🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / database / src / commonMain / kotlin / org / meshtastic / core / database / DatabaseManager.kt
Displaying Raw • Download
core/database/src/commonMain/kotlin/org/meshtastic/core/database/DatabaseManager.kt d5848ad5e02bd9b5b1726bedfa51ceb0faaa0240 (d5848ad5) Text, 98.93 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.database
Tff7b72import T7ee787androidx.datastore.preferences.core.Preferences
Tff7b72import T7ee787androidx.datastore.preferences.core.booleanPreferencesKey
Tff7b72import T7ee787androidx.datastore.preferences.core.edit
Tff7b72import T7ee787androidx.datastore.preferences.core.intPreferencesKey
Tff7b72import T7ee787androidx.datastore.preferences.core.longPreferencesKey
Tff7b72import T7ee787androidx.datastore.preferences.core.stringPreferencesKey
Tff7b72import T7ee787androidx.datastore.preferences.core.stringSetPreferencesKey
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.atomicfu.atomic
Tff7b72import T7ee787kotlinx.atomicfu.locks.SynchronizedObject
Tff7b72import T7ee787kotlinx.atomicfu.locks.synchronized
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CompletableDeferred
Tff7b72import T7ee787kotlinx.coroutines.CoroutineDispatcher
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.CoroutineStart
Tff7b72import T7ee787kotlinx.coroutines.Deferred
Tff7b72import T7ee787kotlinx.coroutines.Dispatchers
Tff7b72import T7ee787kotlinx.coroutines.ExperimentalCoroutinesApi
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.NonCancellable
Tff7b72import T7ee787kotlinx.coroutines.SupervisorJob
Tff7b72import T7ee787kotlinx.coroutines.async
Tff7b72import T7ee787kotlinx.coroutines.currentCoroutineContext
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.ensureActive
Tff7b72import T7ee787kotlinx.coroutines.flow.Flow
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.SharingStarted
Tff7b72import T7ee787kotlinx.coroutines.flow.StateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.catch
Tff7b72import T7ee787kotlinx.coroutines.flow.emitAll
Tff7b72import T7ee787kotlinx.coroutines.flow.first
Tff7b72import T7ee787kotlinx.coroutines.flow.flatMapLatest
Tff7b72import T7ee787kotlinx.coroutines.flow.flow
Tff7b72import T7ee787kotlinx.coroutines.flow.map
Tff7b72import T7ee787kotlinx.coroutines.flow.onEach
Tff7b72import T7ee787kotlinx.coroutines.flow.stateIn
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withContext
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.koin.core.annotation.Single
Tff7b72import T7ee787org.meshtastic.core.common.util.normalizeAddress
Tff7b72import T7ee787org.meshtastic.core.common.util.nowMillis
Tff7b72import T7ee787org.meshtastic.core.database.di.DatabaseDataStore
Tff7b72import T7ee787org.meshtastic.core.di.CoroutineDispatchers
Tff7b72import T7ee787kotlin.concurrent.Volatile
Tff7b72import T7ee787org.meshtastic.core.common.database.DatabaseManager Tff7b72as Te6edf3SharedDatabaseManager
Tff7b72internal Tff7b72const Tff7b72val Te6edf3MAX_CONSECUTIVE_FLOW_POOL_RECOVERIES Tff7b72= T79c0ff3
T8b949e/** Allows two recovered failure streaks while hard-bounding Flow-created detached pools for this manager lifetime. */
Tff7b72internal Tff7b72const Tff7b72val Te6edf3MAX_FLOW_POOL_RECOVERIES_PER_MANAGER_LIFETIME Tff7b72= Te6edf3MAX_CONSECUTIVE_FLOW_POOL_RECOVERIES Tff7b72* T79c0ff2
T8b949e/**
* Hard bound on pools detached by wedge recovery (a `withDb` block that never returned) for this manager lifetime. A
* wedged pool is never closed while an abandoned block may still touch it, so each recovery costs one retained Room
* instance until orderly shutdown.
*/
Tff7b72internal Tff7b72const Tff7b72val Te6edf3MAX_WEDGE_POOL_RECOVERIES_PER_MANAGER_LIFETIME Tff7b72= T79c0ff3
T8b949e/** Returns database names that form either side of an unfinished, crash-recoverable association route. */
Tff7b72internal Tff7b72fun Td2a8ffpendingRouteDbNamesTb4b4b4(Te6edf3preferencesTb4b4b4: Te6edf3PreferencesTb4b4b4)Tb4b4b4: Te6edf3SetTff7b72<Tffa657StringTff7b72> Tff7b72= Te6edf3preferences
Tb4b4b4.Te6edf3asMapTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3asSequenceTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3filter Tb4b4b4{ Tb4b4b4(Te6edf3keyTb4b4b4, Te6edf3_Tb4b4b4) Tff7b72-Tff7b72>
Te6edf3keyTb4b4b4.Te6edf3nameTb4b4b4.Te6edf3startsWithTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3PENDING_SOURCE_DB_FOR_PREFIXTb4b4b4) Tff7b72|Tff7b72|
Te6edf3keyTb4b4b4.Te6edf3nameTb4b4b4.Te6edf3startsWithTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3PENDING_DESTINATION_DB_FOR_PREFIXTb4b4b4)
Tb4b4b4}
Tb4b4b4.Te6edf3mapNotNull Tb4b4b4{ Tb4b4b4(Te6edf3_Tb4b4b4, Te6edf3valueTb4b4b4) Tff7b72-Tff7b72> Te6edf3value Tff7b72as? Tffa657String Tb4b4b4}
Tb4b4b4.Te6edf3toSetTb4b4b4(Tb4b4b4)
T8b949e/** Manages per-device Room database instances for node data, with LRU eviction. */
Tf0883e@SingleTb4b4b4(Te6edf3binds Tff7b72= Tff7b72[Te6edf3DatabaseProviderTff7b72::Te6edf3classTb4b4b4, Te6edf3SharedDatabaseManagerTff7b72::Te6edf3classTff7b72]Tb4b4b4)
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffLargeClassTa5d6ff"Tb4b4b4)
Tf0883e@OptInTb4b4b4(Te6edf3ExperimentalCoroutinesApiTff7b72::Te6edf3classTb4b4b4)
Tff7b72open Tff7b72class T56d364DatabaseManagerTb4b4b4(Tff7b72private Tff7b72val Te6edf3datastoreTb4b4b4: Te6edf3DatabaseDataStoreTb4b4b4, Tff7b72private Tff7b72val Te6edf3dispatchersTb4b4b4: Te6edf3CoroutineDispatchersTb4b4b4) Tb4b4b4:
Te6edf3DatabaseProviderTb4b4b4,
Te6edf3SharedDatabaseManager Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3managerScope Tff7b72= Te6edf3CoroutineScopeTb4b4b4(Te6edf3SupervisorJobTb4b4b4(Tb4b4b4) Tff7b72+ Te6edf3dispatchersTb4b4b4.Te6edf3defaultTb4b4b4)
Tff7b72private Tff7b72val Te6edf3managerJobLock Tff7b72= Te6edf3SynchronizedObjectTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3activeManagerJobs Tff7b72= Te6edf3mutableSetOfTff7b72<Te6edf3JobTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72var Te6edf3managerJobDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
Tff7b72private Tff7b72val Te6edf3mutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3closeMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72enum Tff7b72class T56d364LifecycleState Tb4b4b4{
Te6edf3OPENTb4b4b4,
Te6edf3CLOSINGTb4b4b4,
Te6edf3CLOSEDTb4b4b4,
Tb4b4b4}
Tff7b72private Tff7b72enum Tff7b72class T56d364ReopenOrigin Tb4b4b4{
Te6edf3BOUNDED_OPERATIONTb4b4b4,
Te6edf3WEDGED_OPERATIONTb4b4b4,
Te6edf3FLOW_OBSERVERTb4b4b4,
Tb4b4b4}
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3lifecycleState Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPEN
T8b949e// Per-database bounded-access tracking and per-source write barrier for merges. `withDb` deliberately does NOT
T8b949e// take [mutex] (hot path), so a merge under [mutex] must still drain any in-flight writer that captured the source
T8b949e// DB before folding it away — otherwise a late-committing write is lost when the source is retired. This
T8b949e// dedicated lock (never held across a drain await, so it can't deadlock the merge) tracks live bounded reads and
T8b949e// `withDb` blocks per captured DB instance, plus the writer-admission gate. The gate is armed at the start of an
T8b949e// association attempt: while it is pending, [beginWrite] blocks new writers instead of letting them capture a
T8b949e// DB, so a new `withDb` can never write to `source` once it is being retired, nor land on `dest` before the merge
T8b949e// commits. The gate completes with `source` if the attempt aborts
T8b949e// (drain timeout, cancellation, or pre-commit merge failure) and with `dest` once the merge commits — source is
T8b949e// never restored after the merge commits. The lock is released before any suspend (drain await, gate await, Room
T8b949e// work, merge work, or DataStore work), so it can't deadlock any of them.
Tff7b72private Tff7b72val Te6edf3writerTrackerMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3activeWriters Tff7b72= Te6edf3mutableMapOfTff7b72<Te6edf3MeshtasticDatabaseTb4b4b4, Tffa657IntTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3activeReaders Tff7b72= Te6edf3mutableMapOfTff7b72<Te6edf3MeshtasticDatabaseTb4b4b4, Tffa657IntTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3drainWaiters Tff7b72= Te6edf3mutableMapOfTff7b72<Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3MutableListTff7b72<Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72>Tff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/**
* Per-pool containment lanes for one-shot [withDb] blocks, keyed by Room instance and guarded by
* [writerTrackerMutex]. [beginWrite] hands the lane out in the same critical section that admits the writer, so a
* block always runs on the lane of the pool it was admitted against.
*
* Each lane is a single-parallelism view, which narrows the Room/SQLite churn window without serializing blocks
* across suspension — a suspended callback releases the lane, exactly as the previous process-wide lane behaved.
* Keying by pool is what keeps a callback that blocks its lane thread from stalling every later write: a
* replacement pool published by recovery has its own lane. Entries are dropped once a pool can no longer be
* admitted (reopen detaches it, cached close removes it, [close] clears the map).
*/
Tff7b72private Tff7b72val Te6edf3poolLanes Tff7b72= Te6edf3mutableMapOfTff7b72<Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3CoroutineDispatcherTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/**
* Pools proven wedged that recovery has refused to replace, guarded by [writerTrackerMutex].
*
* Once the lifetime replacement budget is spent, a wedged pool stays published, so without this every later write
* would be admitted against it, wait the full deadline, and abandon another block that can never finish — an
* unbounded pile of parked coroutines and writer registrations under steady ingest. Admission fails fast instead.
* Entries are pruned wherever [poolLanes] entries are: a pool that is replaced or closed is no longer admissible.
*/
Tff7b72private Tff7b72val Te6edf3unrecoverablePools Tff7b72= Te6edf3mutableSetOfTff7b72<Te6edf3MeshtasticDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/**
* Builds the containment lane for one pool. Tests override this because `limitedParallelism` on a test dispatcher
* would replace virtual-time scheduling with a real worker view.
*/
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffcreatePoolLaneTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3CoroutineDispatcher Tff7b72= Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4.Te6edf3limitedParallelismTb4b4b4(T79c0ff1Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3deferredEvictions Tff7b72= Te6edf3mutableSetOfTff7b72<Te6edf3MeshtasticDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72var Te6edf3shutdownWriterDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
T8b949e// Admitted manager-operation tokens. Serialized associateDevice/switch/recovery/eviction/cleanup work and
T8b949e// bounded [withReadDb] callbacks register here before touching a manager-owned pool. Guarded by
T8b949e// [writerTrackerMutex]. [close] arms [managerOperationDrain] and bound-waits for this set to empty BEFORE acquiring
T8b949e// [mutex] for its ownership snapshot, so no admitted operation can resume against a closed pool. New operations are
T8b949e// rejected at admission once lifecycleState != OPEN.
Tff7b72private Tff7b72val Te6edf3activeManagerOperations Tff7b72= Te6edf3mutableSetOfTff7b72<Tffa657AnyTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72var Te6edf3managerOperationDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
T8b949e// Armed at the start of an association attempt; null otherwise. A non-null gate blocks [beginWrite] until the
T8b949e// attempt resolves. It completes with the canonical DB (source on abort, dest on commit) so blocked writers resume
T8b949e// against the right instance. Never read or written outside [writerTrackerMutex].
Tff7b72private Tff7b72var Te6edf3writerGateTb4b4b4: Te6edf3CompletableDeferredTff7b72<Te6edf3MeshtasticDatabaseTff7b72>Tff7b72? Tff7b72= Tff7b72null
Tff7b72private Tff7b72val Te6edf3cacheLimitKey Tff7b72= Te6edf3intPreferencesKeyTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3CACHE_LIMIT_KEYTb4b4b4)
Tff7b72private Tff7b72val Te6edf3legacyCleanedKey Tff7b72= Te6edf3booleanPreferencesKeyTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3LEGACY_DB_CLEANED_KEYTb4b4b4)
Tff7b72private Tff7b72val Te6edf3retiredDbNamesKey Tff7b72= Te6edf3stringSetPreferencesKeyTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3RETIRED_DB_NAMES_KEYTb4b4b4)
Tff7b72private Tff7b72fun Td2a8fflastUsedKeyTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tff7b72= Te6edf3longPreferencesKeyTb4b4b4(Ta5d6ff"Ta5d6ffdb_last_used:Tffd700$Te6edf3dbNameTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffaddrDbKeyTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4) Tff7b72=
Te6edf3stringPreferencesKeyTb4b4b4(Ta5d6ff"Tffd700${Te6edf3DatabaseConstantsTb4b4b4.Te6edf3ADDR_DB_FOR_PREFIXTffd700}Tffd700${Te6edf3normalizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff"Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffpendingSourceDbKeyTb4b4b4(Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4) Tff7b72=
Te6edf3stringPreferencesKeyTb4b4b4(Ta5d6ff"Tffd700${Te6edf3DatabaseConstantsTb4b4b4.Te6edf3PENDING_SOURCE_DB_FOR_PREFIXTffd700}Tffd700${Te6edf3normalizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff"Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffpendingDestinationDbKeyTb4b4b4(Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4) Tff7b72=
Te6edf3stringPreferencesKeyTb4b4b4(Ta5d6ff"Tffd700${Te6edf3DatabaseConstantsTb4b4b4.Te6edf3PENDING_DESTINATION_DB_FOR_PREFIXTffd700}Tffd700${Te6edf3normalizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff"Tb4b4b4)
Tff7b72private Tff7b72var Te6edf3backfillJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null
T8b949e/** Launches and tracks manager-owned work so shutdown waits only for jobs that can still touch owned resources. */
Tff7b72protected Tff7b72fun Td2a8fflaunchManagerWorkTb4b4b4(
Te6edf3dispatcherTb4b4b4: Te6edf3CoroutineDispatcher Tff7b72= Te6edf3dispatchersTb4b4b4.Te6edf3defaultTb4b4b4,
Te6edf3blockTb4b4b4: Te6edf3suspend Te6edf3CoroutineScopeTb4b4b4.Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3Job Tb4b4b4{
Tff7b72lateinit Tff7b72var Te6edf3jobTb4b4b4: Te6edf3Job
Te6edf3job Tff7b72= Te6edf3managerScopeTb4b4b4.Te6edf3launchTb4b4b4(Te6edf3dispatcherTb4b4b4, Te6edf3start Tff7b72= Te6edf3CoroutineStartTb4b4b4.Te6edf3LAZYTb4b4b4, Te6edf3block Tff7b72= Te6edf3blockTb4b4b4)
Tff7b72val Te6edf3accepted Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3managerJobLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72=Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{
Te6edf3activeManagerJobsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3jobTb4b4b4)
T8b949e// A LAZY coroutine can be cancelled after registration but before its body starts. Completion
T8b949e// handlers run in both that case and the normal completion path, so every accepted job releases
T8b949e// its shutdown-drain registration exactly once.
Te6edf3jobTb4b4b4.Te6edf3invokeOnCompletion Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3managerJobLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3activeManagerJobsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3jobTb4b4b4) Tff7b72&Tff7b72& Te6edf3activeManagerJobsTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3managerJobDrainTff7b72?.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72true
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72false
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3acceptedTb4b4b4) Te6edf3jobTb4b4b4.Te6edf3startTb4b4b4(Tb4b4b4) Tff7b72else Te6edf3jobTb4b4b4.Te6edf3cancelTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3job
Tb4b4b4}
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3hasDelayedFirstDeviceBackfill Tff7b72= Tff7b72false
Tff7b72override Tff7b72val Te6edf3cacheLimitTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72> Tff7b72=
Te6edf3datastoreTb4b4b4.Te6edf3data
Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTff7b72[Te6edf3cacheLimitKeyTff7b72] Tff7b72?: Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_CACHE_LIMIT Tb4b4b4}
Tb4b4b4.Te6edf3stateInTb4b4b4(Te6edf3managerScopeTb4b4b4, Te6edf3SharingStartedTb4b4b4.Te6edf3EagerlyTb4b4b4, Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_CACHE_LIMITTb4b4b4)
Tff7b72override Tff7b72fun Td2a8ffgetCurrentCacheLimitTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Int Tff7b72= Te6edf3cacheLimitTb4b4b4.Te6edf3value
Tff7b72override Tff7b72fun Td2a8ffsetCacheLimitTb4b4b4(Te6edf3limitTb4b4b4: Tffa657IntTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3clamped Tff7b72= Te6edf3limitTb4b4b4.Te6edf3coerceInTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3MIN_CACHE_LIMITTb4b4b4, Te6edf3DatabaseConstantsTb4b4b4.Te6edf3MAX_CACHE_LIMITTb4b4b4)
Te6edf3launchManagerWork Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTff7b72[Te6edf3cacheLimitKeyTff7b72] Tff7b72= Te6edf3clamped Tb4b4b4}
T8b949e// Resolve the protected DB only once this deferred work owns the manager mutex. A switch may complete
T8b949e// between scheduling and execution, so capturing currentDbName here could evict the newly active pool.
Te6edf3enforceCacheLimitTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72val Te6edf3dbCache Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657StringTb4b4b4, Te6edf3MeshtasticDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/** Replaced pools that remain live for consumers of an earlier [currentDb] emission until orderly shutdown. */
Tff7b72private Tff7b72val Te6edf3detachedDatabases Tff7b72= Te6edf3mutableListOfTff7b72<Te6edf3NamedDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/** Guarded by [mutex]; counts successful automatic Flow-triggered replacements for this manager's lifetime. */
Tff7b72private Tff7b72var Te6edf3flowPoolRecoveriesThisLifetime Tff7b72= T79c0ff0
T8b949e/** Guarded by [mutex]; counts successful replacements triggered by an abandoned wedged [withDb] block. */
Tff7b72private Tff7b72var Te6edf3wedgePoolRecoveriesThisLifetime Tff7b72= T79c0ff0
T8b949e/** Databases merged and logically retired but kept open — app-wide consumers may still hold references. */
Tff7b72private Tff7b72val Te6edf3logicallyRetired Tff7b72= Te6edf3mutableSetOfTff7b72<Tffa657StringTff7b72>Tb4b4b4(Tb4b4b4)
T8b949e/** Guarded by [mutex]; cleanup is safe only once, before this manager opens a device database. */
Tff7b72private Tff7b72var Te6edf3attemptedStartupRetirementCleanup Tff7b72= Tff7b72false
T8b949e/** Guarded by [mutex]; keeps crash-route recovery from deleting a source after device DBs have been published. */
Tff7b72private Tff7b72var Te6edf3hasOpenedDeviceDatabase Tff7b72= Tff7b72false
Tff7b72private Tff7b72data Tff7b72class T56d364NamedDatabaseTb4b4b4(Tff7b72val Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4, Tff7b72val Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4)
Tff7b72private Tff7b72data Tff7b72class T56d364ShutdownDatabaseTb4b4b4(Tff7b72val Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Tff7b72val Te6edf3dbNamesTb4b4b4: Te6edf3MutableSetTff7b72<Tffa657StringTff7b72>Tb4b4b4)
Tff7b72private Tff7b72data Tff7b72class T56d364ShutdownSnapshotTb4b4b4(Tff7b72val Te6edf3databasesTb4b4b4: Te6edf3ListTff7b72<Te6edf3ShutdownDatabaseTff7b72>Tb4b4b4, Tff7b72val Te6edf3retiredNamesTb4b4b4: Te6edf3ListTff7b72<Tffa657StringTff7b72>Tb4b4b4)
T8b949e/**
* Covers default-pool creation, current-flow publication, and shutdown ownership transfer as one state machine.
* Unlike two independent Kotlin lazy delegates, this lock leaves no gap where shutdown can close a newly built
* default pool before the flow that owns it becomes visible.
*/
Tff7b72private Tff7b72val Te6edf3initializationLock Tff7b72= Te6edf3SynchronizedObjectTb4b4b4(Tb4b4b4)
T8b949e/** Guarded by [initializationLock]. The pool may exist before it is inserted into [dbCache] or published. */
Tff7b72private Tff7b72var Te6edf3initializedDefaultDbTb4b4b4: Te6edf3MeshtasticDatabase? Tff7b72= Tff7b72null
T8b949e/** Guarded by [initializationLock]. Cleared only after shutdown has acquired its contained pool. */
Tff7b72private Tff7b72var Te6edf3currentDbStateTb4b4b4: Te6edf3MutableStateFlowTff7b72<Te6edf3MeshtasticDatabaseTff7b72>Tff7b72? Tff7b72= Tff7b72null
T8b949e/** Caller must hold [initializationLock]. */
Tff7b72private Tff7b72fun Td2a8ffgetOrCreateDefaultDatabaseLockedTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase Tb4b4b4{
Te6edf3initializedDefaultDbTff7b72?.Te6edf3let Tb4b4b4{
Tff7b72return Tffa657it
Tb4b4b4}
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3database Tff7b72= Te6edf3buildDatabaseTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{
Te6edf3runCatching Tb4b4b4{ Te6edf3closeDatabaseTb4b4b4(Te6edf3databaseTb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3onFailure Tb4b4b4{
T8b949e// Retain ownership so the shutdown snapshot can retry instead of losing an open pool.
Te6edf3initializedDefaultDb Tff7b72= Te6edf3database
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to close default database initialized during shutdownTa5d6ff" Tb4b4b4}
Tb4b4b4}
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3initializedDefaultDb Tff7b72= Te6edf3database
Tff7b72return Te6edf3database
Tb4b4b4}
T8b949e/** Builds the default pool once without touching [dbCache]; cache insertion remains manager-mutex-only. */
Tff7b72private Tff7b72fun Td2a8ffgetOrCreateDefaultDatabaseTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{ Te6edf3getOrCreateDefaultDatabaseLockedTb4b4b4(Tb4b4b4) Tb4b4b4}
T8b949e/** Lazily builds and publishes the default pool as one [initializationLock]-guarded ownership transfer. */
Tff7b72private Tff7b72fun Td2a8ffgetOrCreateCurrentDbStateTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3MutableStateFlowTff7b72<Te6edf3MeshtasticDatabaseTff7b72> Tff7b72= Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{
Te6edf3currentDbStateTff7b72?.Te6edf3let Tb4b4b4{
Tff7b72returnTf0883e@synchronized Tffa657it
Tb4b4b4}
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Te6edf3MutableStateFlowTb4b4b4(Te6edf3getOrCreateDefaultDatabaseLockedTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3currentDbState Tff7b72= Tffa657it Tb4b4b4}
Tb4b4b4}
T8b949e/**
* The currently active database. The default DB is opened lazily on first access and every internal publication
* ([switchActiveDatabase], association rollback/release, active-DB reopen recovery) writes [_currentDb] directly,
* so [currentDb].value reflects the new instance on the same program step — no `stateIn`/`filterNotNull` derivation
* that would delay visibility to a coroutine dispatch.
*
* Initialization is deferred until first use so the overridable [buildDatabase] runs only after subclass properties
* are set. It is construction-safe and does not mutate [dbCache]; switch and association paths insert the default
* instance while holding [mutex]. Default creation, flow publication, and shutdown acquisition share
* [initializationLock], so shutdown cannot miss or prematurely close an in-progress publication. Room's `onOpen`
* callback remains lazy until the first query.
*/
Tff7b72private Tff7b72val Te6edf3_currentDbTb4b4b4: Te6edf3MutableStateFlowTff7b72<Te6edf3MeshtasticDatabaseTff7b72>
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3getOrCreateCurrentDbStateTb4b4b4(Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3currentDbTb4b4b4: Te6edf3StateFlowTff7b72<Te6edf3MeshtasticDatabaseTff7b72>
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3_currentDb
T8b949e/**
* Re-latches long-lived DAO flows on database switches and recovers the active Room pool after a reader/writer
* acquisition timeout. A failed query is never replayed on the same pool; publishing the replacement causes
* [flatMapLatest] to start a fresh DAO flow. Concurrent failing collectors converge on the same replacement.
*
* Each collector stops after [MAX_CONSECUTIVE_FLOW_POOL_RECOVERIES] replacements without a successful emission,
* while the manager permits at most [MAX_FLOW_POOL_RECOVERIES_PER_MANAGER_LIFETIME] successful Flow-triggered
* replacements for its lifetime. The second bound prevents intermittent wedge/recover/emission cycles from
* retaining an unbounded number of detached Room instances before orderly shutdown can reclaim them.
*/
Tff7b72override Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffobserveCurrentDbTb4b4b4(Te6edf3queryTb4b4b4: Tb4b4b4(Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3FlowTff7b72<Te6edf3TTff7b72>Tb4b4b4)Tb4b4b4: Te6edf3FlowTff7b72<Te6edf3TTff7b72> Tff7b72= Te6edf3flow Tb4b4b4{
Tff7b72val Te6edf3consecutivePoolRecoveries Tff7b72= Te6edf3atomicTb4b4b4(T79c0ff0Tb4b4b4)
Te6edf3emitAllTb4b4b4(
Te6edf3currentDbTb4b4b4.Te6edf3flatMapLatest Tb4b4b4{ Te6edf3database Tff7b72-Tff7b72>
Te6edf3flow Tb4b4b4{ Te6edf3emitAllTb4b4b4(Te6edf3queryTb4b4b4(Te6edf3databaseTb4b4b4)Tb4b4b4.Te6edf3onEach Tb4b4b4{ Te6edf3consecutivePoolRecoveriesTb4b4b4.Te6edf3value Tff7b72= T79c0ff0 Tb4b4b4}Tb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3catch Tb4b4b4{ Te6edf3failure Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Te6edf3failure Tff7b72is Te6edf3CancellationExceptionTb4b4b4) Tff7b72throw Te6edf3failure
Tff7b72val Te6edf3exception Tff7b72= Te6edf3failure Tff7b72as? Te6edf3Exception Tff7b72?: Tff7b72throw Te6edf3failure
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isDbPoolAcquireTimeoutExceptionTb4b4b4(Te6edf3exceptionTb4b4b4)Tb4b4b4) Tff7b72throw Te6edf3failure
Tff7b72if Tb4b4b4(
Te6edf3consecutivePoolRecoveriesTb4b4b4.Te6edf3value Tff7b72>Tff7b72= Te6edf3MAX_CONSECUTIVE_FLOW_POOL_RECOVERIES Tff7b72&Tff7b72&
Te6edf3shouldRethrowFlowFailureTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3exception
Tb4b4b4}
Te6edf3consecutivePoolRecoveriesTb4b4b4.Te6edf3incrementAndGetTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3reopened Tff7b72=
Tff7b72try Tb4b4b4{
T8b949e// Once a timeout is classified, finish or reject the manager-level ownership transfer
T8b949e// atomically even if this collector is cancelled. A cancellable reopen could otherwise
T8b949e// build a replacement and abandon the cache/publication handoff halfway through.
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3reopenFlowDatabaseIfStillCurrentTb4b4b4(Te6edf3databaseTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4) Te6edf3recoveryFailureTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3exceptionTb4b4b4.Te6edf3addSuppressedTb4b4b4(Te6edf3recoveryFailureTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3recoveryFailureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to recover active DB after a Flow pool timeoutTa5d6ff" Tb4b4b4}
Tff7b72throw Te6edf3exception
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3reopened Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffReopened active DB after a Room Flow connection-pool timeoutTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72else Tff7b72if Tb4b4b4(Te6edf3shouldRethrowFlowFailureTb4b4b4(Te6edf3databaseTb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3exception
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
T8b949e/** Checks flow ownership without lazily rebuilding [currentDb] while shutdown is clearing its publication. */
Tff7b72private Tff7b72fun Td2a8ffshouldRethrowFlowFailureTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3isStillCurrent Tff7b72= Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{ Te6edf3currentDbStateTff7b72?.Te6edf3value Tff7b72=Tff7b72=Tff7b72= Te6edf3database Tb4b4b4}
Tff7b72return Te6edf3isStillCurrent Tff7b72|Tff7b72| Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPEN
Tb4b4b4}
Tff7b72private Tff7b72val Te6edf3_currentAddress Tff7b72= Te6edf3MutableStateFlowTff7b72<Tffa657String?Tff7b72>Tb4b4b4(Tff7b72nullTb4b4b4)
Tff7b72val Te6edf3currentAddressTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657String?Tff7b72> Tff7b72= Te6edf3_currentAddress
T8b949e/**
* Name of the currently active database. Tracked explicitly rather than recomputed from the address, because
* cross-transport aliasing ([associateDevice]) decouples the two: a secondary transport's address maps to the DB
* claimed by the first transport, which `buildDbName(address)` would never produce. Written under [mutex].
*/
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3currentDbNameTb4b4b4: Tffa657String Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAME
T8b949e/** Initialize the active database for [address]. */
Tff7b72suspend Tff7b72fun Td2a8ffinitTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4) Tb4b4b4{
Te6edf3switchActiveDatabaseTb4b4b4(Te6edf3addressTb4b4b4)
Tb4b4b4}
T8b949e/** Returns a cached database or builds one. Every caller must hold [mutex]. */
Tff7b72private Tff7b72fun Td2a8ffgetOrOpenDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase Tff7b72= Te6edf3dbCacheTb4b4b4.Te6edf3getOrPutTb4b4b4(Te6edf3dbNameTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3dbName Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Te6edf3getOrCreateDefaultDatabaseTb4b4b4(Tb4b4b4) Tff7b72else Te6edf3buildDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffcheckOpenTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3lifecycleState Tff7b72=Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffDatabaseManager is closing or closedTa5d6ff" Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Admits one operation that may touch manager-owned database pools and drains it on completion.
*
* Registration happens before the operation captures a pool or waits on [mutex] (after the lifecycle check); the
* token is removed in `finally` so cancellation, failure, and timeout all release admission. A [close] that
* observes a non-empty admission set arms [managerOperationDrain] and bound-waits on it before taking its ownership
* snapshot, so an admitted callback cannot resume against a closed pool.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffwithManagerOperationTb4b4b4(Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Te6edf3TTb4b4b4)Tb4b4b4: Te6edf3T Tb4b4b4{
Tff7b72val Te6edf3token Tff7b72= Tffa657AnyTb4b4b4(Tb4b4b4)
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Te6edf3activeManagerOperationsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3tokenTb4b4b4)
Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72return Te6edf3blockTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3activeManagerOperationsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3tokenTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3activeManagerOperationsTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3managerOperationDrainTff7b72?.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Builds a new [MeshtasticDatabase] for [dbName]. Tests override this to control file placement (temp directory
* instead of the platform data dir). Production delegates to the platform-specific [getDatabaseBuilder].
*/
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffbuildDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase Tff7b72= Te6edf3getDatabaseBuilderTb4b4b4(Te6edf3dbNameTb4b4b4)Tb4b4b4.Te6edf3buildTb4b4b4(Tb4b4b4)
T8b949e/**
* Resolves the DB name to use for [address], honoring a cross-transport alias when one exists. A secondary
* transport (e.g. TCP) that has been unified with a node points at the DB the first transport (e.g. BLE) claimed;
* without an alias this falls back to the address-hashed name — today's default — for a first-time or primary
* connection. See [associateDevice].
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffresolveDbNameTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4, Te6edf3canReclaimRecoveredSourceTb4b4b4: Tffa657BooleanTb4b4b4)Tb4b4b4: Tffa657String Tb4b4b4{
Tff7b72val Te6edf3fallback Tff7b72= Te6edf3buildDbNameTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3fallback Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Tff7b72return Te6edf3fallback
Tff7b72val Te6edf3transportAddress Tff7b72= Te6edf3address Tff7b72?: Tff7b72return Te6edf3fallback
Tff7b72val Te6edf3aliasKey Tff7b72= Te6edf3addrDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)
Tff7b72val Te6edf3pendingSourceKey Tff7b72= Te6edf3pendingSourceDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)
Tff7b72val Te6edf3pendingDestinationKey Tff7b72= Te6edf3pendingDestinationDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)
Tff7b72val Te6edf3prefs Tff7b72= Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3pendingSource Tff7b72= Te6edf3prefsTff7b72[Te6edf3pendingSourceKeyTff7b72]
Tff7b72val Te6edf3pendingDestination Tff7b72= Te6edf3prefsTff7b72[Te6edf3pendingDestinationKeyTff7b72]
Tff7b72if Tb4b4b4(Te6edf3pendingSource Tff7b72=Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3pendingDestination Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tff7b72return Te6edf3prefsTff7b72[Te6edf3aliasKeyTff7b72] Tff7b72?: Te6edf3fallback
Tff7b72if Tb4b4b4(Te6edf3pendingSource Tff7b72=Tff7b72= Tff7b72null Tff7b72|Tff7b72| Te6edf3pendingDestination Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceKeyTb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationKeyTb4b4b4)
Tb4b4b4}
Tff7b72return Te6edf3prefsTff7b72[Te6edf3aliasKeyTff7b72] Tff7b72?: Te6edf3fallback
Tb4b4b4}
T8b949e// A pending route is only intent. The destination's merge marker is the durable proof that the data copy
T8b949e// committed. Verify it before publishing either database so no caller can write to a merged-away fallback.
Tff7b72val Te6edf3destinationDb Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3getOrOpenDatabaseTb4b4b4(Te6edf3pendingDestinationTb4b4b4) Tb4b4b4}
Tff7b72val Te6edf3mergeCommitted Tff7b72= Te6edf3verifyMergeMarkerTb4b4b4(Te6edf3destinationDbTb4b4b4, Te6edf3pendingSourceTb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3mergeCommittedTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceKeyTb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationKeyTb4b4b4)
Tb4b4b4}
Tff7b72return Te6edf3prefsTff7b72[Te6edf3aliasKeyTff7b72] Tff7b72?: Te6edf3fallback
Tb4b4b4}
T8b949e// If this process may already have published the source pool, its durable retirement and in-memory
T8b949e// protection are one cancellation-atomic transition. Otherwise cancellation after DataStore commits but
T8b949e// before logicallyRetired is updated could let cache eviction close/delete a pool still held by consumers.
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Tffa657itTff7b72[Te6edf3aliasKeyTff7b72] Tff7b72= Te6edf3pendingDestination
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceKeyTb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationKeyTb4b4b4)
Tffa657itTff7b72[Te6edf3lastUsedKeyTb4b4b4(Te6edf3pendingDestinationTb4b4b4)Tff7b72] Tff7b72= Te6edf3nowMillis
Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72] Tff7b72= Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4) Tff7b72+ Te6edf3pendingSource
Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3canReclaimRecoveredSourceTb4b4b4) Te6edf3logicallyRetiredTb4b4b4.Te6edf3addTb4b4b4(Te6edf3pendingSourceTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3canReclaimRecoveredSourceTb4b4b4) Tb4b4b4{
Te6edf3physicallyRetireDatabaseTb4b4b4(Te6edf3pendingSourceTb4b4b4)
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{
Ta5d6ff"Ta5d6ffRepaired pending route from Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3pendingSourceTb4b4b4)Tffd700}Ta5d6ff to Ta5d6ff" Tff7b72+ Te6edf3anonymizeDbNameTb4b4b4(Te6edf3pendingDestinationTb4b4b4)
Tb4b4b4}
Tff7b72return Te6edf3pendingDestination
Tb4b4b4}
T8b949e/** Reads the destination merge marker used as commit proof for pending-route recovery. */
Tff7b72protected Tff7b72open Tff7b72suspend Tff7b72fun Td2a8ffverifyMergeMarkerTb4b4b4(Te6edf3destinationTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3sourceNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3destinationTb4b4b4.Te6edf3mergeMarkerDaoTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3isMergedTb4b4b4(Te6edf3sourceNameTb4b4b4)
T8b949e/** Switch active database to the one associated with [address]. Serialized via mutex. */
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffswitchActiveDatabaseTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4) Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Te6edf3cleanupPersistedRetirementsAtStartupTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3dbName Tff7b72= Te6edf3resolveDbNameTb4b4b4(Te6edf3addressTb4b4b4, Te6edf3canReclaimRecoveredSource Tff7b72= Tff7b72!Te6edf3hasOpenedDeviceDatabaseTb4b4b4)
T8b949e// Remember the previously active DB name (any) so we can record its last-used time as well.
Tff7b72val Te6edf3previousDbName Tff7b72= Te6edf3currentDbName
T8b949e// Fast path: no-op only when both the selected address and its resolved database are already active.
T8b949e// resolveDbName() may repair a committed pending route to a different destination for the same address.
Tff7b72if Tb4b4b4(Te6edf3_currentAddressTb4b4b4.Te6edf3value Tff7b72=Tff7b72= Te6edf3address Tff7b72&Tff7b72& Te6edf3currentDbName Tff7b72=Tff7b72= Te6edf3dbNameTb4b4b4) Tb4b4b4{
Te6edf3markLastUsedTb4b4b4(Te6edf3dbNameTb4b4b4)
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// Build/open Room DB off the main thread
Tff7b72val Te6edf3db Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3getOrOpenDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3dbName Tff7b72!Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Te6edf3hasOpenedDeviceDatabase Tff7b72= Tff7b72true
T8b949e// Emit the new DB BEFORE closing the old ones. flatMapLatest collectors on
T8b949e// currentDb will cancel their in-flight queries on the previous database once
T8b949e// the new value is emitted. Closing the old pool first would race with those
T8b949e// collectors, causing "Connection pool is closed" crashes.
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72= Te6edf3db
Te6edf3currentDbName Tff7b72= Te6edf3dbName
Tb4b4b4}
Te6edf3_currentAddressTb4b4b4.Te6edf3value Tff7b72= Te6edf3address
Te6edf3markLastUsedTb4b4b4(Te6edf3dbNameTb4b4b4)
T8b949e// Also mark the previous DB as used "just now" so LRU has an accurate, recent timestamp
Te6edf3markLastUsedTb4b4b4(Te6edf3previousDbNameTb4b4b4)
T8b949e// Do NOT close the previous DB synchronously here. Even though _currentDb has been
T8b949e// updated, in-flight `withDb` calls may still hold a reference to the old database
T8b949e// (captured before the emission). Closing the connection pool while those queries are
T8b949e// executing causes "Connection pool is closed" crashes. Instead, let LRU eviction
T8b949e// (enforceCacheLimit) handle cleanup — it only runs on databases that are not the
T8b949e// active target and have not been used recently.
Te6edf3schedulePostSwitchMaintenanceTb4b4b4(Te6edf3dbName Tff7b72= Te6edf3dbNameTb4b4b4, Te6edf3db Tff7b72= Te6edf3dbTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffSwitched active DB to Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff for address Tffd700${Te6edf3anonymizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Schedules deferred maintenance that runs after switching the active database. Posts work to [managerScope] on
* [dispatchers.io] so the switch path is not blocked by filesystem or search-index I/O.
*
* In production this schedules LRU cache-limit enforcement, legacy-DB cleanup, and FTS search-index backfill.
* In-memory test fixtures override it to no-op because they do not have a filesystem-backed database directory and
* must not access platform context singletons (e.g. `ContextServices.app`).
*/
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffschedulePostSwitchMaintenanceTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4, Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
T8b949e// Defer LRU eviction so switch is not blocked by filesystem work
Te6edf3launchManagerWorkTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3enforceCacheLimitTb4b4b4(Tb4b4b4) Tb4b4b4}
T8b949e// One-time cleanup: remove legacy DB if present and not active
Te6edf3launchManagerWorkTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3cleanupLegacyDbIfNeededTb4b4b4(Te6edf3activeDbName Tff7b72= Te6edf3dbNameTb4b4b4) Tb4b4b4}
T8b949e// Backfill FTS search index for any text messages missing messageText.
T8b949e// On the first real device DB, defer this so it does not starve the single DB connection while
T8b949e// the UI is collecting startup flows. The default DB should not consume the cold-start delay.
Tff7b72val Te6edf3shouldDelayBackfill Tff7b72= Te6edf3dbName Tff7b72!Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAME Tff7b72&Tff7b72& Tff7b72!Te6edf3hasDelayedFirstDeviceBackfill
Tff7b72if Tb4b4b4(Te6edf3shouldDelayBackfillTb4b4b4) Te6edf3hasDelayedFirstDeviceBackfill Tff7b72= Tff7b72true
Te6edf3scheduleSearchIndexBackfillTb4b4b4(Te6edf3dbName Tff7b72= Te6edf3dbNameTb4b4b4, Te6edf3db Tff7b72= Te6edf3dbTb4b4b4, Te6edf3shouldDelayBackfill Tff7b72= Te6edf3shouldDelayBackfillTb4b4b4)
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffCyclomaticComplexMethodTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffLongMethodTa5d6ff"Tb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffassociateDeviceTb4b4b4(
Te6edf3addressTb4b4b4: Tffa657StringTb4b4b4,
Te6edf3nodeNumTb4b4b4: Tffa657IntTb4b4b4,
Te6edf3deviceIdTb4b4b4: Tffa657String?Tb4b4b4,
Te6edf3isSessionActiveTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657BooleanTb4b4b4,
Tb4b4b4) Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72fun Td2a8ffensureAssociationActiveTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isSessionActiveTb4b4b4(Tb4b4b4)Tb4b4b4) Tff7b72throw Te6edf3StaleAssociationExceptionTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3_currentAddressTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3addressTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{
Ta5d6ff"Ta5d6ffIgnored stale database association for Tffd700${Te6edf3anonymizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff; active transport is Ta5d6ff" Tff7b72+
Te6edf3anonymizeAddressTb4b4b4(Te6edf3_currentAddressTb4b4b4.Te6edf3valueTb4b4b4)
Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Tff7b72val Te6edf3sourceName Tff7b72= Te6edf3currentDbName
T8b949e// Never claim or merge into the sentinel "no device" DB.
Tff7b72if Tb4b4b4(Te6edf3sourceName Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Tff7b72returnTf0883e@withLock
T8b949e// The device-id claim is the durable one (node numbers renumber under firmware 2.8); the
T8b949e// node-num claim stays as the fallback for hardware without a device id, for lockdown
T8b949e// sessions (device_id zeroed), and for claims written by older app versions. Writes always
T8b949e// refresh both keys so either lookup path resolves on the next connection.
Tff7b72val Te6edf3deviceKey Tff7b72= Te6edf3validDeviceIdOrNullTb4b4b4(Te6edf3deviceIdTb4b4b4)Tff7b72?.Te6edf3letTb4b4b4(Tff7b72::Te6edf3deviceDbPrefKeyTb4b4b4)
Tff7b72val Te6edf3nodeKey Tff7b72= Te6edf3nodeDbPrefKeyTb4b4b4(Te6edf3nodeNumTb4b4b4)
Tff7b72val Te6edf3prefs Tff7b72= Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3claimed Tff7b72= Te6edf3resolveDbClaimTb4b4b4(Te6edf3prefsTb4b4b4, Te6edf3deviceKeyTb4b4b4, Te6edf3nodeKeyTb4b4b4)
Tff7b72suspend Tff7b72fun Td2a8ffwriteClaimsTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tff7b72= Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Te6edf3deviceKeyTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3key Tff7b72-Tff7b72> Tffa657itTff7b72[Te6edf3keyTff7b72] Tff7b72= Te6edf3dbName Tb4b4b4}
Tffa657itTff7b72[Te6edf3nodeKeyTff7b72] Tff7b72= Te6edf3dbName
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72when Tb4b4b4{
Te6edf3claimed Tff7b72=Tff7b72= Tff7b72null Tff7b72-Tff7b72> Tb4b4b4{
T8b949e// First transport to learn this device: its current DB becomes the device's canonical DB.
T8b949e// No address alias is needed — a primary connection already resolves to this DB via
T8b949e// buildDbName.
Te6edf3writeClaimsTb4b4b4(Te6edf3sourceNameTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffClaimed Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3sourceNameTb4b4b4)Tffd700}Ta5d6ff as canonical DB for node Tffd700$Te6edf3nodeNumTa5d6ff" Tb4b4b4}
Tb4b4b4}
Te6edf3claimed Tff7b72=Tff7b72= Te6edf3sourceName Tff7b72-Tff7b72> Tb4b4b4{
T8b949e// Already unified — backfill or refresh any stale/missing routing metadata atomically.
T8b949e// This also repairs a post-merge routing failure (merge committed but DataStore edit failed):
T8b949e// the next connect reaches this branch and writes claims + alias in one edit without
T8b949e// re-copying source (the merge marker prevents duplicate data).
Tff7b72val Te6edf3addressKey Tff7b72= Te6edf3addrDbKeyTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72val Te6edf3needsDeviceKey Tff7b72= Te6edf3deviceKey Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3prefsTff7b72[Te6edf3deviceKeyTff7b72] Tff7b72!Tff7b72= Te6edf3sourceName
Tff7b72val Te6edf3needsNodeKey Tff7b72= Te6edf3prefsTff7b72[Te6edf3nodeKeyTff7b72] Tff7b72!Tff7b72= Te6edf3sourceName
Tff7b72val Te6edf3needsAlias Tff7b72= Te6edf3prefsTff7b72[Te6edf3addressKeyTff7b72] Tff7b72!Tff7b72= Te6edf3sourceName
Tff7b72val Te6edf3pendingSourceKey Tff7b72= Te6edf3pendingSourceDbKeyTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72val Te6edf3pendingDestinationKey Tff7b72= Te6edf3pendingDestinationDbKeyTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72val Te6edf3needsPendingCleanup Tff7b72=
Te6edf3prefsTff7b72[Te6edf3pendingSourceKeyTff7b72] Tff7b72!Tff7b72= Tff7b72null Tff7b72|Tff7b72| Te6edf3prefsTff7b72[Te6edf3pendingDestinationKeyTff7b72] Tff7b72!Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3needsDeviceKey Tff7b72|Tff7b72| Te6edf3needsNodeKey Tff7b72|Tff7b72| Te6edf3needsAlias Tff7b72|Tff7b72| Te6edf3needsPendingCleanupTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Te6edf3deviceKeyTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3key Tff7b72-Tff7b72> Tff7b72if Tb4b4b4(Te6edf3needsDeviceKeyTb4b4b4) Tffa657itTff7b72[Te6edf3keyTff7b72] Tff7b72= Te6edf3sourceName Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3needsNodeKeyTb4b4b4) Tffa657itTff7b72[Te6edf3nodeKeyTff7b72] Tff7b72= Te6edf3sourceName
Tff7b72if Tb4b4b4(Te6edf3needsAliasTb4b4b4) Tffa657itTff7b72[Te6edf3addressKeyTff7b72] Tff7b72= Te6edf3sourceName
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceKeyTb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationKeyTb4b4b4)
Tffa657itTff7b72[Te6edf3lastUsedKeyTb4b4b4(Te6edf3sourceNameTb4b4b4)Tff7b72] Tff7b72= Te6edf3nowMillis
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffRefreshed routing metadata for Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3sourceNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72else Tff7b72-Tff7b72> Tb4b4b4{
T8b949e// Secondary transport reached an already-known node: fold this DB into the canonical one,
T8b949e// switch the active DB to it, alias this address to it, and retire the now-merged source.
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3source Tff7b72= Te6edf3_currentDbTb4b4b4.Te6edf3value
Tff7b72val Te6edf3dest Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3getOrOpenDatabaseTb4b4b4(Te6edf3claimedTb4b4b4) Tb4b4b4}
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
T8b949e// Arm the writer-admission gate before any drain or merge work. The complete armed lifetime is
T8b949e// enclosed by the try/finally below, so every CancellationException, Exception, and Error
T8b949e// releases blocked writers onto source before commit or destination after commit.
Tff7b72val Te6edf3gate Tff7b72= Te6edf3CompletableDeferredTff7b72<Te6edf3MeshtasticDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3transportAddress Tff7b72= Te6edf3address
Tff7b72var Te6edf3mergeCommitted Tff7b72= Tff7b72false
Tff7b72var Te6edf3retirementPersisted Tff7b72= Tff7b72false
T8b949e// Publishes and completes this exact gate once. A repeated cleanup call is harmless and cannot
T8b949e// clear a later association's gate.
Tff7b72suspend Tff7b72fun Td2a8ffreleaseWriterGateTb4b4b4(Te6edf3canonicalDbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3canonicalNameTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3released Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3writerGate Tff7b72!Tff7b72=Tff7b72= Te6edf3gateTb4b4b4) Tb4b4b4{
Tff7b72false
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3writerGate Tff7b72= Tff7b72null
Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72= Te6edf3canonicalDb
Te6edf3currentDbName Tff7b72= Te6edf3canonicalName
Tff7b72true
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3releasedTb4b4b4) Te6edf3gateTb4b4b4.Te6edf3completeTb4b4b4(Te6edf3canonicalDbTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72suspend Tff7b72fun Td2a8ffclearPendingRouteBestEffortTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3cleanupFailureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3cleanupFailureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to clear aborted pending database routeTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72try Tb4b4b4{
T8b949e// Start the outer try before arming the gate. Once writerGate is assigned, no Throwable can
T8b949e// escape without running the identity-checked release in finally.
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3checkTb4b4b4(Te6edf3writerGate Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffDatabase writer gate already armedTa5d6ff" Tb4b4b4}
Te6edf3writerGate Tff7b72= Te6edf3gate
Tb4b4b4}
Tff7b72try Tb4b4b4{
T8b949e// Phase 1: drain every writer admitted against source before this gate was armed.
Tff7b72val Te6edf3drained Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3drainWritersTb4b4b4(Te6edf3sourceTb4b4b4, Te6edf3sourceNameTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3drainedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffAborted merge of Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3sourceNameTb4b4b4)Tffd700}Ta5d6ff into Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3claimedTb4b4b4)Tffd700}Ta5d6ff: writer drain timed out; kept Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3sourceNameTb4b4b4)Tffd700}Ta5d6ff activeTa5d6ff"
Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
T8b949e// Phase 2: persist address-scoped intent, commit the database merge and marker, then
T8b949e// finalize every route and remove intent in one DataStore transaction. If finalization
T8b949e// fails after the merge commits, a later switch verifies the marker and repairs the
T8b949e// alias before publishing.
Te6edf3withContextTb4b4b4(Te6edf3NonCancellable Tff7b72+ Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tffa657itTff7b72[Te6edf3pendingSourceDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tff7b72] Tff7b72= Te6edf3sourceName
Tffa657itTff7b72[Te6edf3pendingDestinationDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tff7b72] Tff7b72= Te6edf3claimed
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3mergeDatabasesTb4b4b4(Te6edf3sourceTb4b4b4, Te6edf3destTb4b4b4, Te6edf3sourceNameTb4b4b4, Te6edf3isSessionActiveTb4b4b4)
Te6edf3mergeCommitted Tff7b72= Tff7b72true
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Te6edf3deviceKeyTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3key Tff7b72-Tff7b72> Tffa657itTff7b72[Te6edf3keyTff7b72] Tff7b72= Te6edf3claimed Tb4b4b4}
Tffa657itTff7b72[Te6edf3nodeKeyTff7b72] Tff7b72= Te6edf3claimed
Tffa657itTff7b72[Te6edf3addrDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tff7b72] Tff7b72= Te6edf3claimed
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingSourceDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tb4b4b4)
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3pendingDestinationDbKeyTb4b4b4(Te6edf3transportAddressTb4b4b4)Tb4b4b4)
Tffa657itTff7b72[Te6edf3lastUsedKeyTb4b4b4(Te6edf3claimedTb4b4b4)Tff7b72] Tff7b72= Te6edf3nowMillis
Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72] Tff7b72= Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4) Tff7b72+ Te6edf3sourceName
Te6edf3ensureAssociationActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3retirementPersisted Tff7b72= Tff7b72true
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{
Ta5d6ff"Ta5d6ffUnified Tffd700${Te6edf3anonymizeDbNameTb4b4b4(
Te6edf3sourceNameTb4b4b4,
Tb4b4b4)Tffd700}Ta5d6ff into Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3claimedTb4b4b4)Tffd700}Ta5d6ff for node Tffd700$Te6edf3nodeNumTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3mergeCommittedTb4b4b4) Tb4b4b4{
Te6edf3clearPendingRouteBestEffortTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72when Tb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Tff7b72is Te6edf3CancellationExceptionTb4b4b4,
Tff7b72is Te6edf3StaleAssociationExceptionTb4b4b4,
Tff7b72-Tff7b72> Tff7b72throw Te6edf3failure
Tff7b72is Te6edf3Exception Tff7b72-Tff7b72> Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3mergeCommittedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffMerge into Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3claimedTb4b4b4)Tffd700}Ta5d6ff failed; Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffkept Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3sourceNameTb4b4b4)Tffd700}Ta5d6ff activeTa5d6ff"
Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffRouting metadata for Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3claimedTb4b4b4)Tffd700}Ta5d6ff failed after merge Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffcommit; destination remains active and the pending route will Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffrepair the address alias on a later switchTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Tff7b72else Tff7b72-Tff7b72> Tff7b72throw Te6edf3failure
Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3mergeCommittedTb4b4b4) Tb4b4b4{
Te6edf3releaseWriterGateTb4b4b4(Te6edf3destTb4b4b4, Te6edf3claimedTb4b4b4)
Te6edf3recordLogicalRetirementTb4b4b4(Te6edf3sourceNameTb4b4b4, Te6edf3persistIntent Tff7b72= Tff7b72!Te6edf3retirementPersistedTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3releaseWriterGateTb4b4b4(Te6edf3sourceTb4b4b4, Te6edf3sourceNameTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3_Tb4b4b4: Te6edf3StaleAssociationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{
Ta5d6ff"Ta5d6ffAborted stale database association for Tffd700${Te6edf3anonymizeAddressTb4b4b4(Te6edf3addressTb4b4b4)Tffd700}Ta5d6ff after transport-session rolloverTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Logically retires a database whose contents have been merged into another.
*
* Physical close/delete is deferred — the merged source was published through [currentDb]; app-wide Flow, Paging,
* UI, worker, and one-shot read consumers may still hold its Room instance. Physically closing it now can surface
* "Connection pool is closed" to those readers (see [switchActiveDatabase] and [reopenActiveDatabaseIfStillCurrent]
* no-sync-close discipline). Retirement intent is persisted so the next manager lifetime can safely remove the file
* before opening a device DB; [close] also performs orderly physical teardown when the platform invokes it.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffrecordLogicalRetirementTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4, Te6edf3persistIntentTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
T8b949e// associateDevice already owns the manager mutex. Reacquiring it here would deadlock finalization.
Te6edf3logicallyRetiredTb4b4b4.Te6edf3addTb4b4b4(Te6edf3dbNameTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3persistIntentTb4b4b4) Te6edf3persistRetirementIntentTb4b4b4(Te6edf3dbNameTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffLogically retired merged DB Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff; physical cleanup deferredTa5d6ff" Tb4b4b4}
Tb4b4b4}
T8b949e/** Persists merge retirement separately as a fallback when post-commit route finalization failed. */
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffpersistRetirementIntentTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72] Tff7b72= Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4) Tff7b72+ Te6edf3dbName Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to persist retirement for Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Reclaims retirement intents left by a previous application process. This runs at most once and before opening any
* device DB, which is the only point where no consumer from this process can hold a retired Room instance.
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffcleanupPersistedRetirementsAtStartupTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3attemptedStartupRetirementCleanupTb4b4b4) Tff7b72return
Te6edf3attemptedStartupRetirementCleanup Tff7b72= Tff7b72true
Tff7b72val Te6edf3retiredNames Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3cancellationTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3cancellation
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to read persisted database retirementsTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Te6edf3retiredNamesTb4b4b4.Te6edf3forEach Tb4b4b4{
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Te6edf3physicallyRetireDatabaseTb4b4b4(Tffa657itTb4b4b4)
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/** Deletes one retired file, then atomically clears its retirement and last-used metadata. */
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffphysicallyRetireDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3deleteDatabaseFilesTb4b4b4(Te6edf3dbNameTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to delete retired database Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Tff7b72try Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{
Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3lastUsedKeyTb4b4b4(Te6edf3dbNameTb4b4b4)Tb4b4b4)
Tff7b72val Te6edf3remaining Tff7b72= Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3dbName
Tff7b72if Tb4b4b4(Te6edf3remainingTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3retiredDbNamesKeyTb4b4b4) Tff7b72else Tffa657itTff7b72[Te6edf3retiredDbNamesKeyTff7b72] Tff7b72= Te6edf3remaining
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffPhysically retired merged DB Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
T8b949e// The file deletion is idempotent. Retain the intent so a later process retries metadata cleanup.
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to clear retirement metadata for Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Closes and removes a cached database by name. Safe to call even if the database was already closed or not in the
* cache. Does NOT delete the underlying file — the database can be re-opened on next access.
*
* Room KMP is configured with a single-connection pool on every platform and has no common auto-close timeout, so
* an idle cached database keeps that connection open until explicitly closed. This method is the primary mechanism
* for releasing it when a database is no longer the active target. The caller must hold [mutex].
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72protected Tff7b72open Tff7b72suspend Tff7b72fun Td2a8ffcloseCachedDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3database Tff7b72= Te6edf3dbCacheTff7b72[Te6edf3dbNameTff7b72] Tff7b72?: Tff7b72return
Tff7b72try Tb4b4b4{
Te6edf3closeDatabaseTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to close cached database Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72throw Te6edf3failure
Tb4b4b4}
Te6edf3dbCacheTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dbNameTb4b4b4)
T8b949e// A closed pool can never be admitted again, so its lane and wedge verdict are dead weight.
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3poolLanesTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3databaseTb4b4b4)
Te6edf3unrecoverablePoolsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffClosed inactive database Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff to free connectionsTa5d6ff" Tb4b4b4}
Tb4b4b4}
T8b949e/** Room-close seam used by deterministic shutdown tests. */
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffcloseDatabaseTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72= Te6edf3databaseTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
T8b949e/**
* Reopens the active database under [mutex], but only if it hasn't switched since the caller snapshotted it.
*
* The replaced Room instance is intentionally left open for the rest of the process. [currentDb] reads [_currentDb]
* directly, so every publication is visible to app-wide collectors on the same program step — but there is no
* deterministic handoff point where every collector has stopped using the previous instance.
*
* Returns the reopened DB, or null if another coroutine switched databases or shutdown has started.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffreopenFlowDatabaseIfStillCurrentTb4b4b4(Te6edf3expectedDbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase? Tb4b4b4{
Tff7b72val Te6edf3expectedDbName Tff7b72=
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPEN Tff7b72|Tff7b72| Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72!Tff7b72=Tff7b72= Te6edf3expectedDbTb4b4b4) Tff7b72return Tff7b72null
Te6edf3currentDbName
Tb4b4b4}
Tff7b72return Te6edf3reopenActiveDatabaseIfStillCurrentTb4b4b4(Te6edf3expectedDbTb4b4b4, Te6edf3expectedDbNameTb4b4b4, Te6edf3ReopenOriginTb4b4b4.Te6edf3FLOW_OBSERVERTb4b4b4)
Tb4b4b4}
T8b949e/**
* Reports whether automatic replacement is still allowed for [origin]. Each automatic recovery retains one detached
* Room instance until orderly shutdown, so an intermittent wedge/recover cycle must not run unbounded. Caller must
* hold [mutex].
*/
Tff7b72private Tff7b72fun Td2a8ffhasReachedRecoveryLimitTb4b4b4(Te6edf3originTb4b4b4: Te6edf3ReopenOriginTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Tb4b4b4(Te6edf3countTb4b4b4, Te6edf3limitTb4b4b4) Tff7b72=
Tff7b72when Tb4b4b4(Te6edf3originTb4b4b4) Tb4b4b4{
Te6edf3ReopenOriginTb4b4b4.Te6edf3FLOW_OBSERVER Tff7b72-Tff7b72>
Te6edf3flowPoolRecoveriesThisLifetime Te6edf3to Te6edf3MAX_FLOW_POOL_RECOVERIES_PER_MANAGER_LIFETIME
Te6edf3ReopenOriginTb4b4b4.Te6edf3WEDGED_OPERATION Tff7b72-Tff7b72>
Te6edf3wedgePoolRecoveriesThisLifetime Te6edf3to Te6edf3MAX_WEDGE_POOL_RECOVERIES_PER_MANAGER_LIFETIME
T8b949e// Failure-triggered recovery is already bounded by the failing call itself.
Te6edf3ReopenOriginTb4b4b4.Te6edf3BOUNDED_OPERATION Tff7b72-Tff7b72> Tff7b72return Tff7b72false
Tb4b4b4}
Tff7b72return Tb4b4b4(Te6edf3count Tff7b72>Tff7b72= Te6edf3limitTb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3reached Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Te6edf3reachedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffDB recovery limit reached for Tffd700$Te6edf3originTa5d6ff (Tffd700$Te6edf3countTa5d6ff); refusing another automatic replacementTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Counts one successful automatic replacement against [origin]'s lifetime budget. Caller must hold [mutex]. */
Tff7b72private Tff7b72fun Td2a8ffrecordRecoveryTb4b4b4(Te6edf3originTb4b4b4: Te6edf3ReopenOriginTb4b4b4) Tb4b4b4{
Tff7b72when Tb4b4b4(Te6edf3originTb4b4b4) Tb4b4b4{
Te6edf3ReopenOriginTb4b4b4.Te6edf3FLOW_OBSERVER Tff7b72-Tff7b72> Te6edf3flowPoolRecoveriesThisLifetime Tff7b72+Tff7b72= T79c0ff1
Te6edf3ReopenOriginTb4b4b4.Te6edf3WEDGED_OPERATION Tff7b72-Tff7b72> Te6edf3wedgePoolRecoveriesThisLifetime Tff7b72+Tff7b72= T79c0ff1
Te6edf3ReopenOriginTb4b4b4.Te6edf3BOUNDED_OPERATION Tff7b72-Tff7b72> Tffa657Unit
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffreopenActiveDatabaseIfStillCurrentTb4b4b4(
Te6edf3expectedDbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4,
Te6edf3expectedDbNameTb4b4b4: Tffa657StringTb4b4b4,
Te6edf3originTb4b4b4: Te6edf3ReopenOriginTb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase? Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tff7b72returnTf0883e@withManagerOperation Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72!Tff7b72=Tff7b72= Te6edf3expectedDb Tff7b72|Tff7b72| Te6edf3currentDbName Tff7b72!Tff7b72= Te6edf3expectedDbNameTb4b4b4) Tff7b72returnTf0883e@withManagerOperation Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3hasReachedRecoveryLimitTb4b4b4(Te6edf3originTb4b4b4)Tb4b4b4) Tff7b72returnTf0883e@withManagerOperation Tff7b72null
Tff7b72val Te6edf3registered Tff7b72=
Te6edf3dbCacheTff7b72[Te6edf3expectedDbNameTff7b72]
Tff7b72?: Tff7b72if Tb4b4b4(Te6edf3expectedDbName Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{ Te6edf3initializedDefaultDb Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72null
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3registered Tff7b72!Tff7b72=Tff7b72= Te6edf3expectedDbTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffDB recovery: active registration changed before reopen; skipping active DB reopenTa5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withManagerOperation Tff7b72null
Tb4b4b4}
T8b949e// Build a fresh instance directly (not through getOrPut) before touching the cache,
T8b949e// so a failed or cancelled build leaves the existing cache entry and _currentDb consistent.
Tff7b72val Te6edf3reopened Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3buildDatabaseTb4b4b4(Te6edf3expectedDbNameTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{
Te6edf3runCatching Tb4b4b4{ Te6edf3closeDatabaseTb4b4b4(Te6edf3reopenedTb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3onFailure Tb4b4b4{
Te6edf3detachedDatabasesTb4b4b4.Te6edf3addTb4b4b4(Te6edf3NamedDatabaseTb4b4b4(Te6edf3expectedDbNameTb4b4b4, Te6edf3reopenedTb4b4b4)Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to close database built during shutdown; retained for shutdown retryTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72returnTf0883e@withManagerOperation Tff7b72null
Tb4b4b4}
Te6edf3dbCacheTff7b72[Te6edf3expectedDbNameTff7b72] Tff7b72= Te6edf3reopened
Tff7b72if Tb4b4b4(Te6edf3expectedDbName Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4) Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{ Te6edf3initializedDefaultDb Tff7b72= Te6edf3reopened Tb4b4b4}
Tb4b4b4}
Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72= Te6edf3reopened
Te6edf3recordRecoveryTb4b4b4(Te6edf3originTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3detachedDatabasesTb4b4b4.Te6edf3none Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3database Tff7b72=Tff7b72=Tff7b72= Te6edf3expectedDb Tb4b4b4}Tb4b4b4) Tb4b4b4{
Te6edf3detachedDatabasesTb4b4b4.Te6edf3addTb4b4b4(Te6edf3NamedDatabaseTb4b4b4(Te6edf3expectedDbNameTb4b4b4, Te6edf3expectedDbTb4b4b4)Tb4b4b4)
Tb4b4b4}
T8b949e// The replaced instance can no longer be admitted (beginWrite only ever hands out the published pool), so
T8b949e// drop its lane and any wedge verdict. An abandoned block keeps running on the lane it already captured.
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3poolLanesTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3expectedDbTb4b4b4)
Te6edf3unrecoverablePoolsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3expectedDbTb4b4b4)
Tb4b4b4}
T8b949e// Intentionally do not close expectedDb here. The public currentDb Flow exposes _currentDb directly,
T8b949e// so downstream flatMapLatest collectors may still be using the replaced Room instance after this
T8b949e// function emits the reopened DB. Closing the old pool here can surface "Connection pool is closed"
T8b949e// to app-wide DB observers that do not have closed-pool recovery. This mirrors switchActiveDatabase's
T8b949e// no-sync-close discipline. [close] owns the detached-pool set and reclaims every replaced instance after
T8b949e// all
T8b949e// application consumers and admitted writers have stopped.
Te6edf3reopened
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Deadline for one [withDb] call. Tests override this: a non-positive value waits without a bound, which is the
* only way a fixture can park a writer across a virtual-time idle period without the scheduler auto-advancing into
* this deadline.
*/
Tff7b72protected Tff7b72open Tff7b72val Te6edf3withDbTimeoutMillisTb4b4b4: Tffa657Long
Tff7b72getTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3WITH_DB_TIMEOUT_MS
T8b949e/**
* Executes [block] once against the admitted current DB instance, bounding the caller's wait at
* [withDbTimeoutMillis].
*
* A callback is never replayed after it starts: an arbitrary block can perform one side effect and then fail, so
* transparently invoking it again against another pool could duplicate or split a logical write. Pool-timeout
* recovery may reopen the active database for future calls, but the original failure is still propagated.
*
* Containment is per pool, not process-wide: the callback runs on the pool's containment lane ([poolLanes]) in a
* manager-scoped child coroutine that owns the writer registration and releases it in its own `finally`. Only the
* *wait* is bounded — a started block still runs to completion under [NonCancellable], so a bounded DB-critical
* section is never torn apart mid-write. A block that never returns (Room 3.x logs and retries connection-pool
* acquisition instead of throwing, so a leaked permit hangs forever) is abandoned rather than cancelled: the caller
* fails with [DatabaseOperationTimeoutException] and the active pool is reopened, so later calls are admitted
* against the replacement pool and its fresh lane instead of queueing behind the wedge.
*
* Long-lived Flow/Paging reads must stay out of `withDb`; see [observeCurrentDb] and [withReadDb].
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffwithDbTb4b4b4(Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3TTb4b4b4)Tb4b4b4: Te6edf3T? Tb4b4b4{
Tff7b72val Te6edf3queuedAt Tff7b72= Te6edf3nowMillis
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3admission Tff7b72= Te6edf3beginWriteTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3blockStarted Tff7b72= Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3execution Tff7b72= Te6edf3launchDbBlockTb4b4b4(Te6edf3admissionTb4b4b4, Te6edf3blockStartedTb4b4b4, Te6edf3blockTb4b4b4)
Tff7b72val Te6edf3timeoutMillis Tff7b72= Te6edf3withDbTimeoutMillis
Tff7b72val Te6edf3completed Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3timeoutMillis Tff7b72<Tff7b72= T79c0ff0Tb4b4b4) Te6edf3executionTb4b4b4.Te6edf3joinTb4b4b4(Tb4b4b4) Tff7b72else Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3timeoutMillisTb4b4b4) Tb4b4b4{ Te6edf3executionTb4b4b4.Te6edf3joinTb4b4b4(Tb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3completed Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Te6edf3abandonWedgedDbBlockTb4b4b4(Te6edf3admissionTb4b4b4, Te6edf3blockStartedTb4b4b4, Te6edf3executionTb4b4b4, Te6edf3timeoutMillisTb4b4b4)
T8b949e// Re-check cancellation so a stale caller does not continue after the DB releases.
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3elapsedMillis Tff7b72= Te6edf3nowMillis Tff7b72- Te6edf3queuedAt
Tff7b72if Tb4b4b4(Te6edf3elapsedMillis Tff7b72>Tff7b72= Te6edf3WITH_DB_SLOW_OPERATION_MSTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb took Tffd700${Te6edf3elapsedMillisTffd700}Ta5d6ffms including lane wait; persistent slow logs indicate the DB access Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffpath should be revisitedTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tff7b72val Te6edf3failure Tff7b72= Te6edf3executionTb4b4b4.Te6edf3getCompletionExceptionOrNullTb4b4b4(Tb4b4b4) Tff7b72?: Tff7b72return Te6edf3executionTb4b4b4.Te6edf3getCompletedTb4b4b4(Tb4b4b4)
Te6edf3handleDbBlockFailureTb4b4b4(Te6edf3admissionTb4b4b4, Te6edf3failureTb4b4b4)
Tb4b4b4}
T8b949e/**
* Starts [block] on the admitted pool's lane.
*
* The coroutine owns the writer registration and releases it in its own `finally`, so an abandoned block can never
* double-release it and always deregisters if it eventually returns. It starts [CoroutineStart.ATOMIC] so that
* `finally` also runs when the coroutine is cancelled before its first dispatch, and it completes [blockStarted]
* only once it is about to invoke the callback — a call still waiting for its lane has performed no side effect and
* is aborted at the cancellation check instead of running late.
*/
Tff7b72private Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8fflaunchDbBlockTb4b4b4(
Te6edf3admissionTb4b4b4: Te6edf3AdmittedDatabaseTb4b4b4,
Te6edf3blockStartedTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4,
Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3TTb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3DeferredTff7b72<Te6edf3TTff7b72> Tff7b72= Te6edf3managerScopeTb4b4b4.Te6edf3asyncTb4b4b4(Te6edf3admissionTb4b4b4.Te6edf3laneTb4b4b4, Te6edf3start Tff7b72= Te6edf3CoroutineStartTb4b4b4.Te6edf3ATOMICTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3db Tff7b72= Te6edf3admissionTb4b4b4.Te6edf3database
Tff7b72try Tb4b4b4{
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Te6edf3blockStartedTb4b4b4.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4)
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3blockTb4b4b4(Te6edf3dbTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3endWriteTb4b4b4(Te6edf3dbTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Stops admitting work against a pool that stayed published after recovery declined to replace it.
*
* Recovery declines once the lifetime replacement budget is spent, and it can also fail outright. Either way the
* wedged instance remains [_currentDb], so every later admission would park another block on it forever. Marking it
* turns those calls into an immediate failure. Recovery during shutdown is not a wedge verdict, and neither is a
* pool that has already been replaced, so both are left unmarked.
*/
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffquarantineWedgedPoolIfStillPublishedTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tff7b72return
Tff7b72val Te6edf3quarantined Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPEN Tff7b72|Tff7b72| Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72!Tff7b72=Tff7b72= Te6edf3databaseTb4b4b4) Tb4b4b4{
Tff7b72false
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3unrecoverablePoolsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3quarantinedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffMarked the active DB pool unrecoverable after wedge recovery was refused; further database Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffoperations fail fast until the active database is replacedTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Resolves a [withDb] call whose bounded wait expired: cancels a block that never started, or abandons a started
* one and recovers the pool.
*
* A started block is not cancellable by design, so it keeps its lane and writer registration until it returns — a
* merge drain or shutdown drain that waits on it still resolves on its own bound and logs. Reopening the active
* pool is what unwedges the app: the wedged instance moves to the detached set (protected from eviction, reclaimed
* by [close]) and later callers admit against the replacement pool and its fresh lane.
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffThrowsCountTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffabandonWedgedDbBlockTb4b4b4(
Te6edf3admissionTb4b4b4: Te6edf3AdmittedDatabaseTb4b4b4,
Te6edf3blockStartedTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4,
Te6edf3executionTb4b4b4: Te6edf3DeferredTff7b72<Tff7b72*Tff7b72>Tb4b4b4,
Te6edf3timeoutMillisTb4b4b4: Tffa657LongTb4b4b4,
Tb4b4b4)Tb4b4b4: Te6edf3Nothing Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3blockStartedTb4b4b4.Te6edf3isCompletedTb4b4b4) Tb4b4b4{
Te6edf3executionTb4b4b4.Te6edf3cancelTb4b4b4(Tb4b4b4)
T8b949e// Re-check: the block may have acquired the lane between the timeout and the cancellation, in which case it
T8b949e// is inside NonCancellable and must be abandoned instead of reported as never started.
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3blockStartedTb4b4b4.Te6edf3isCompletedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffwithDb waited Tffd700${Te6edf3timeoutMillisTffd700}Ta5d6ffms for the DB lane; its callback never startedTa5d6ff" Tb4b4b4}
Tff7b72throw Te6edf3DatabaseOperationTimeoutExceptionTb4b4b4(
Ta5d6ff"Ta5d6ffwithDb did not start within Tffd700${Te6edf3timeoutMillisTffd700}Ta5d6ffms; an earlier block is wedged on this databaseTa5d6ff"Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72val Te6edf3reopened Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3reopenActiveDatabaseIfStillCurrentTb4b4b4(Te6edf3admissionTb4b4b4.Te6edf3databaseTb4b4b4, Te6edf3admissionTb4b4b4.Te6edf3nameTb4b4b4, Te6edf3ReopenOriginTb4b4b4.Te6edf3WEDGED_OPERATIONTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3recoveryCancelTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3recoveryCancel
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3recoveryFailureTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3recoveryFailureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffwithDb: failed to reopen active DB after abandoning a wedged callbackTa5d6ff" Tb4b4b4}
Tff7b72null
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3reopened Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Te6edf3quarantineWedgedPoolIfStillPublishedTb4b4b4(Te6edf3admissionTb4b4b4.Te6edf3databaseTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3reopened Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb callback exceeded Tffd700${Te6edf3timeoutMillisTffd700}Ta5d6ffms; abandoned it and reopened the active DB so later Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffcalls run on a fresh connection poolTa5d6ff"
Tb4b4b4} Tff7b72else Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb callback exceeded Tffd700${Te6edf3timeoutMillisTffd700}Ta5d6ffms; abandoned it but the active DB was not reopenedTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tff7b72throw Te6edf3DatabaseOperationTimeoutExceptionTb4b4b4(
Ta5d6ff"Ta5d6ffwithDb callback did not finish within Tffd700${Te6edf3timeoutMillisTffd700}Ta5d6ffms and was abandoned without being replayedTa5d6ff"Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
T8b949e/**
* Executes one bounded read without writer admission or the serialized write-containment lane. Active database
* publication is synchronous and the captured pool is registered against eviction until the callback completes. The
* read is also admitted into the shutdown drain, is never replayed automatically, and new reads are rejected once
* shutdown begins.
*/
Tff7b72override Tff7b72suspend Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffwithReadDbTb4b4b4(Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3TTb4b4b4)Tb4b4b4: Te6edf3T Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Tff7b72val Te6edf3database Tff7b72= Te6edf3beginReadTb4b4b4(Tb4b4b4)
Tff7b72try Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3result Tff7b72= Te6edf3blockTb4b4b4(Te6edf3databaseTb4b4b4)
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Te6edf3result
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3endReadTb4b4b4(Te6edf3databaseTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Captures and registers the currently published pool so eviction cannot close it while the callback is active. */
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffbeginReadTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3MeshtasticDatabase Tff7b72= Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Te6edf3_currentDbTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3database Tff7b72-Tff7b72> Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72= Tb4b4b4(Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72?: T79c0ff0Tb4b4b4) Tff7b72+ T79c0ff1 Tb4b4b4}
Tb4b4b4}
T8b949e/** Releases a bounded reader and retries an eviction that was deferred while the captured pool was in use. */
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffendReadTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3retryEviction Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3remaining Tff7b72= Tb4b4b4(Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72?: T79c0ff1Tb4b4b4) Tff7b72- T79c0ff1
Tff7b72if Tb4b4b4(Te6edf3remaining Tff7b72<Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3activeReadersTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72= Te6edf3remaining
Tb4b4b4}
Tff7b72!Te6edf3hasActiveDatabaseAccessLockedTb4b4b4(Te6edf3databaseTb4b4b4) Tff7b72&Tff7b72& Te6edf3deferredEvictionsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3databaseTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3retryEviction Tff7b72&Tff7b72& Te6edf3lifecycleState Tff7b72=Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{
Te6edf3launchManagerWorkTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3enforceCacheLimitTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Caller must hold [writerTrackerMutex]. */
Tff7b72private Tff7b72fun Td2a8ffhasActiveDatabaseAccessLockedTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Tb4b4b4(Te6edf3activeWritersTff7b72[Te6edf3databaseTff7b72] Tff7b72?: T79c0ff0Tb4b4b4) Tff7b72> T79c0ff0 Tff7b72|Tff7b72| Tb4b4b4(Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72?: T79c0ff0Tb4b4b4) Tff7b72> T79c0ff0
T8b949e/**
* Atomically snapshots the canonical active DB (held by [_currentDb], which initializes lazily to the default DB)
* and registers a writer against it.
*
* If an association attempt is in flight, the writer-admission gate is armed. The caller snapshots that gate under
* [writerTrackerMutex], awaits it outside the lock, then retries admission from the beginning. Selecting the active
* database and registering the writer happen in the same critical section, so an association cannot arm its gate
* between those operations. This guarantees a new `withDb` never writes to a DB that is being retired, nor lands on
* `dest` before its data exists.
*/
Tff7b72private Tff7b72data Tff7b72class T56d364AdmittedDatabaseTb4b4b4(
Tff7b72val Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4,
Tff7b72val Te6edf3nameTb4b4b4: Tffa657StringTb4b4b4,
T8b949e/** Containment lane of [database]; see [poolLanes]. */
Tff7b72val Te6edf3laneTb4b4b4: Te6edf3CoroutineDispatcherTb4b4b4,
Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffbeginWriteTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3AdmittedDatabase Tb4b4b4{
Tff7b72while Tb4b4b4(Tff7b72trueTb4b4b4) Tb4b4b4{
Tff7b72var Te6edf3admittedTb4b4b4: Te6edf3AdmittedDatabase? Tff7b72= Tff7b72null
Tff7b72val Te6edf3gate Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3checkOpenTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3pendingGate Tff7b72= Te6edf3writerGate
Tff7b72if Tb4b4b4(Te6edf3pendingGate Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3db Tff7b72= Te6edf3_currentDbTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(Te6edf3db Tff7b72in Te6edf3unrecoverablePoolsTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3DatabaseOperationTimeoutExceptionTb4b4b4(
Ta5d6ff"Ta5d6ffdatabase pool is wedged and recovery is exhausted; refusing to start another Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffoperation that cannot finishTa5d6ff"Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Te6edf3activeWritersTff7b72[Te6edf3dbTff7b72] Tff7b72= Tb4b4b4(Te6edf3activeWritersTff7b72[Te6edf3dbTff7b72] Tff7b72?: T79c0ff0Tb4b4b4) Tff7b72+ T79c0ff1
Te6edf3admitted Tff7b72=
Te6edf3AdmittedDatabaseTb4b4b4(
Te6edf3database Tff7b72= Te6edf3dbTb4b4b4,
Te6edf3name Tff7b72= Te6edf3currentDbNameTb4b4b4,
Te6edf3lane Tff7b72= Te6edf3poolLanesTb4b4b4.Te6edf3getOrPutTb4b4b4(Te6edf3dbTb4b4b4) Tb4b4b4{ Te6edf3createPoolLaneTb4b4b4(Tb4b4b4) Tb4b4b4}Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Te6edf3pendingGate
Tb4b4b4}
Te6edf3admittedTff7b72?.Te6edf3let Tb4b4b4{
Tff7b72return Tffa657it
Tb4b4b4}
Tff7b72val Te6edf3pendingGate Tff7b72= Te6edf3checkNotNullTb4b4b4(Te6edf3gateTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffWriter admission produced neither a database nor a gateTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3released Tff7b72=
Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3WRITER_GATE_TIMEOUT_MSTb4b4b4) Tb4b4b4{
Te6edf3pendingGateTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4)
Tff7b72true
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3released Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3IllegalStateExceptionTb4b4b4(
Ta5d6ff"Ta5d6ffTimed out waiting Tffd700${Te6edf3WRITER_GATE_TIMEOUT_MSTffd700}Ta5d6ffms for database writer admission gateTa5d6ff"Tb4b4b4,
Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Deregisters a writer and releases any merge waiting for [db] to quiesce. Cancellation-safe (see call site). */
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffendWriteTb4b4b4(Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3retryEviction Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3remaining Tff7b72= Tb4b4b4(Te6edf3activeWritersTff7b72[Te6edf3dbTff7b72] Tff7b72?: T79c0ff1Tb4b4b4) Tff7b72- T79c0ff1
Tff7b72val Te6edf3drained Tff7b72= Te6edf3remaining Tff7b72<Tff7b72= T79c0ff0
Tff7b72if Tb4b4b4(Te6edf3drainedTb4b4b4) Tb4b4b4{
Te6edf3activeWritersTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dbTb4b4b4)
Te6edf3drainWaitersTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dbTb4b4b4)Tff7b72?.Te6edf3forEach Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3activeWritersTff7b72[Te6edf3dbTff7b72] Tff7b72= Te6edf3remaining
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3activeWritersTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Te6edf3shutdownWriterDrainTff7b72?.Te6edf3completeTb4b4b4(Tffa657UnitTb4b4b4)
Tff7b72!Te6edf3hasActiveDatabaseAccessLockedTb4b4b4(Te6edf3dbTb4b4b4) Tff7b72&Tff7b72& Te6edf3deferredEvictionsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dbTb4b4b4)
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3retryEviction Tff7b72&Tff7b72& Te6edf3lifecycleState Tff7b72=Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tb4b4b4{
Te6edf3launchManagerWorkTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3enforceCacheLimitTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Folds [source] into [dest]. Override in tests to inject merge failures. Production delegates to
* [DatabaseMerger.merge]; the merge runs in a single transaction so a crash rolls back cleanly and the destination
* is never left half-merged.
*/
Tff7b72protected Tff7b72open Tff7b72suspend Tff7b72fun Td2a8ffmergeDatabasesTb4b4b4(
Te6edf3sourceTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4,
Te6edf3destTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4,
Te6edf3sourceNameTb4b4b4: Tffa657StringTb4b4b4,
Te6edf3isAssociationActiveTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657BooleanTb4b4b4,
Tb4b4b4) Tb4b4b4{
Te6edf3DatabaseMergerTb4b4b4.Te6edf3mergeTb4b4b4(Te6edf3sourceTb4b4b4, Te6edf3destTb4b4b4, Te6edf3sourceNameTb4b4b4, Te6edf3isAssociationActiveTb4b4b4)
Tb4b4b4}
T8b949e/**
* Test-only snapshot of the writer tracker: total live writers and total pending drain waiters. Both are zero once
* every association attempt has released its gate and drained its source — a non-zero pair after a quiescent period
* indicates a leaked writer or waiter.
*/
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffdebugWriterCountsTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3PairTff7b72<Tffa657IntTb4b4b4, Tffa657IntTff7b72> Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3activeWritersTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3sumTb4b4b4(Tb4b4b4) Te6edf3to Te6edf3drainWaitersTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3sumOf Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3size Tb4b4b4} Tb4b4b4}
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffdebugReaderCountTb4b4b4(Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabase? Tff7b72= Tff7b72nullTb4b4b4)Tb4b4b4: Tffa657Int Tff7b72= Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3database Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Te6edf3activeReadersTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3sumTb4b4b4(Tb4b4b4) Tff7b72else Te6edf3activeReadersTff7b72[Te6edf3databaseTff7b72] Tff7b72?: T79c0ff0
Tb4b4b4}
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffdebugEnforceCacheLimitTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3enforceCacheLimitTb4b4b4(Tb4b4b4)
T8b949e/** Test-only visibility for asserting an association did not leak writer admission. */
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffdebugWriterGateArmedTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3writerGate Tff7b72!Tff7b72= Tff7b72null Tb4b4b4}
T8b949e/** Test-only visibility for cancellation-atomic pending-route recovery assertions. */
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffdebugIsLogicallyRetiredTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3dbName Tff7b72in Te6edf3logicallyRetired Tb4b4b4}
T8b949e/** Test-only visibility for deterministic shutdown assertions. */
Tff7b72internal Tff7b72fun Td2a8ffdebugAcceptingWritesTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3lifecycleState Tff7b72=Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPEN
T8b949e/**
* Suspends until every writer that captured [db] before this call has finished, so a merge never snapshots [db]
* while a write is still in flight (and then loses it when [db] is retired). Bounded by [WRITER_DRAIN_TIMEOUT_MS]
* so a wedged writer can't pin the merge — and [mutex] — forever.
*
* Returns `true` if all writers drained (or none were active), `false` on timeout. The caller must abort the merge
* and roll the active DB back to source on `false`.
*
* The waiter is removed in a [finally] block on every exit path — success, timeout, and external cancellation — so
* a stale [CompletableDeferred] never leaks into [drainWaiters]. The cleanup runs under [NonCancellable] so
* cancellation during cleanup doesn't skip the removal.
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffdrainWritersTb4b4b4(Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3waiter Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Tb4b4b4(Te6edf3activeWritersTff7b72[Te6edf3dbTff7b72] Tff7b72?: T79c0ff0Tb4b4b4) Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tff7b72return Tff7b72true
Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3drainWaitersTb4b4b4.Te6edf3getOrPutTb4b4b4(Te6edf3dbTb4b4b4) Tb4b4b4{ Te6edf3mutableListOfTb4b4b4(Tb4b4b4) Tb4b4b4}Tb4b4b4.Te6edf3addTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3drained Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3WRITER_DRAIN_TIMEOUT_MSTb4b4b4) Tb4b4b4{ Te6edf3waiterTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3drained Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTimed out draining writers on Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff before mergeTa5d6ff" Tb4b4b4}
Tff7b72return Tff7b72false
Tb4b4b4}
Tff7b72return Tff7b72true
Tb4b4b4} Tff7b72finally Tb4b4b4{
T8b949e// Remove our waiter on every exit path. On success, endWrite may have already removed the
T8b949e// entire list — the removal is idempotent. On timeout or cancellation, the waiter is still
T8b949e// registered and must be cleaned up so a late endWrite doesn't complete a dead deferred.
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3list Tff7b72= Te6edf3drainWaitersTff7b72[Te6edf3dbTff7b72]
Tff7b72if Tb4b4b4(Te6edf3list Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3listTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3waiterTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3listTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Te6edf3drainWaitersTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dbTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Maps a failed [withDb] execution onto the existing recovery policy and rethrows.
*
* The failed callback is never replayed; recovery only reopens the pool so *future* calls succeed.
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffReturnCountTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffThrowsCountTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffCyclomaticComplexMethodTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffhandleDbBlockFailureTb4b4b4(Te6edf3admissionTb4b4b4: Te6edf3AdmittedDatabaseTb4b4b4, Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4)Tb4b4b4: Te6edf3Nothing Tb4b4b4{
Tff7b72val Te6edf3db Tff7b72= Te6edf3admissionTb4b4b4.Te6edf3database
Tff7b72val Te6edf3active Tff7b72= Te6edf3admissionTb4b4b4.Te6edf3name
Tff7b72if Tb4b4b4(Te6edf3failure Tff7b72is Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
T8b949e// The execution coroutine was cancelled from outside this call — only manager shutdown does that. The
T8b949e// caller itself is still active (its own cancellation surfaces from join), so do not impersonate it.
Tff7b72throw Te6edf3IllegalStateExceptionTb4b4b4(Ta5d6ff"Ta5d6ffwithDb work was cancelled before completingTa5d6ff"Tb4b4b4, Te6edf3failureTb4b4b4)
Tb4b4b4}
Tff7b72val Te6edf3e Tff7b72= Te6edf3failure Tff7b72as? Te6edf3Exception Tff7b72?: Tff7b72throw Te6edf3failure
T8b949e// Shutdown in progress: do not touch _currentDb (its getter calls checkOpen()) nor attempt a reopen that
T8b949e// would build a pool. Propagate the original failure with its message intact.
Tff7b72if Tb4b4b4(Te6edf3lifecycleState Tff7b72!Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4) Tff7b72throw Te6edf3e
Tff7b72val Te6edf3currentDb Tff7b72= Te6edf3_currentDbTb4b4b4.Te6edf3value
Tff7b72if Tb4b4b4(Te6edf3currentDb Tff7b72!Tff7b72=Tff7b72= Te6edf3db Tff7b72&Tff7b72& Te6edf3isDbClosedExceptionTb4b4b4(Te6edf3eTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb: database closed during switch (Tffd700${Te6edf3eTb4b4b4.Te6edf3messageTffd700}Ta5d6ff); callback will not be replayed automaticallyTa5d6ff"
Tb4b4b4}
Tff7b72throw Te6edf3e
Tb4b4b4}
T8b949e// Same active DB but Room's connection pool is wedged. Reopen for future calls, but do not replay this
T8b949e// callback: it may already have completed an earlier side effect before the timeout surfaced.
Tff7b72if Tb4b4b4(Te6edf3currentDb Tff7b72=Tff7b72=Tff7b72= Te6edf3db Tff7b72&Tff7b72& Te6edf3isDbPoolAcquireTimeoutExceptionTb4b4b4(Te6edf3eTb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3reopened Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3reopenActiveDatabaseIfStillCurrentTb4b4b4(Te6edf3dbTb4b4b4, Te6edf3activeTb4b4b4, Te6edf3ReopenOriginTb4b4b4.Te6edf3BOUNDED_OPERATIONTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3recoveryCancelTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3recoveryCancel
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3recoveryFailureTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3eTb4b4b4.Te6edf3addSuppressedTb4b4b4(Te6edf3recoveryFailureTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3recoveryFailureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffwithDb: failed to reopen active DB after a connection-pool timeoutTa5d6ff" Tb4b4b4}
Tff7b72null
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3reopened Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb: reopened active DB after transient Room connection-pool timeout; Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6fffailed callback was not replayedTa5d6ff"
Tb4b4b4} Tff7b72else Tb4b4b4{
Ta5d6ff"Ta5d6ffwithDb: active DB was not reopened during timeout recovery; failed callback was not replayedTa5d6ff"
Tb4b4b4}
Tb4b4b4}
Tff7b72throw Te6edf3e
Tb4b4b4}
Tff7b72throw Te6edf3e
Tb4b4b4}
T8b949e/**
* Runs one bounded DB-critical block under the caller's own writer admission, preserved through cancellation.
*
* Unlike [withDb] this has no lane and no wait bound: the only caller is manager-owned search-index backfill, which
* runs in a cancellable manager job and can legitimately take much longer than [WITH_DB_TIMEOUT_MS].
*/
Tff7b72private Tff7b72suspend Tff7b72fun Tff7b72<Te6edf3TTff7b72> Td2a8ffrunCancellableDbBlockTb4b4b4(Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3blockTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3MeshtasticDatabaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3TTb4b4b4)Tb4b4b4: Te6edf3T Tb4b4b4{
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3result Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3blockTb4b4b4(Te6edf3dbTb4b4b4) Tb4b4b4}
Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3ensureActiveTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3result
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffisDbClosedExceptionTb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3isDbPoolAcquireTimeoutExceptionTb4b4b4(Te6edf3eTb4b4b4) Tff7b72|Tff7b72|
Te6edf3generateSequenceTff7b72<Te6edf3ThrowableTff7b72>Tb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cause Tb4b4b4}
Tb4b4b4.Te6edf3any Tb4b4b4{ Te6edf3throwable Tff7b72-Tff7b72>
Tff7b72val Te6edf3msg Tff7b72= Te6edf3throwableTb4b4b4.Te6edf3messageTff7b72?.Te6edf3lowercaseTb4b4b4(Tb4b4b4) Tff7b72?: Tff7b72returnTf0883e@any Tff7b72false
Tff7b72val Te6edf3hasDbContext Tff7b72= Te6edf3DB_TERMSTb4b4b4.Te6edf3any Tb4b4b4{ Tffa657it Tff7b72in Te6edf3msg Tb4b4b4}
Tb4b4b4(Ta5d6ff"Ta5d6ffclosedTa5d6ff" Tff7b72in Te6edf3msg Tff7b72&Tff7b72& Te6edf3hasDbContextTb4b4b4) Tff7b72|Tff7b72| Ta5d6ff"Ta5d6ffdatabase is lockedTa5d6ff" Tff7b72in Te6edf3msg Tff7b72|Tff7b72| Ta5d6ff"Ta5d6ffsqlite_busyTa5d6ff" Tff7b72in Te6edf3msg
Tb4b4b4}
Tff7b72internal Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72private Tff7b72const Tff7b72val Te6edf3BACKFILL_COLD_START_DELAY_MS Tff7b72= T79c0ff2Te6edf3_000L
Tff7b72private Tff7b72const Tff7b72val Te6edf3WITH_DB_SLOW_OPERATION_MS Tff7b72= T79c0ff1Te6edf3_000L
T8b949e/**
* Upper bound on the lane wait plus the callback itself. Writer admission is bounded separately by
* [WRITER_GATE_TIMEOUT_MS] and runs before this deadline starts, so a call blocked behind an association can
* take the sum of the two. Generous enough that no legitimate one-shot write reaches it, and short enough that
* a pool wedge surfaces as a failure with recovery instead of a hang.
*/
Tff7b72internal Tff7b72const Tff7b72val Te6edf3WITH_DB_TIMEOUT_MS Tff7b72= T79c0ff3T79c0ff0Te6edf3_000L
T8b949e/**
* Upper bound on how long a merge waits for in-flight writers on the source DB to drain (see [drainWriters]).
*/
Tff7b72private Tff7b72const Tff7b72val Te6edf3WRITER_DRAIN_TIMEOUT_MS Tff7b72= T79c0ff5Te6edf3_000L
Tff7b72private Tff7b72const Tff7b72val Te6edf3WRITER_GATE_TIMEOUT_MS Tff7b72= T79c0ff3T79c0ff0Te6edf3_000L
Tff7b72val Te6edf3DB_TERMS Tff7b72= Te6edf3listOfTb4b4b4(Ta5d6ff"Ta5d6ffpoolTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffdatabaseTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffconnectionTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffsqliteTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72const Tff7b72val Te6edf3ROOM_POOL_ACQUIRE_TIMEOUT_PHRASE Tff7b72= Ta5d6ff"Ta5d6fftimed out attempting to acquireTa5d6ff"
Tff7b72private Tff7b72const Tff7b72val Te6edf3ROOM_READER_CONNECTION_PHRASE Tff7b72= Ta5d6ff"Ta5d6ffreader connectionTa5d6ff"
Tff7b72private Tff7b72const Tff7b72val Te6edf3ROOM_WRITER_CONNECTION_PHRASE Tff7b72= Ta5d6ff"Ta5d6ffwriter connectionTa5d6ff"
T8b949e/**
* Room KMP currently exposes pool-acquire timeouts as exception message text instead of a stable common typed
* signal. Keep this fallback narrow so BLE/GATT/transport connection errors do not trigger DB reopen recovery.
*/
Tff7b72private Tff7b72fun Td2a8ffisRoomPoolAcquireTimeoutMessageTb4b4b4(Te6edf3messageTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72=
Te6edf3ROOM_POOL_ACQUIRE_TIMEOUT_PHRASE Tff7b72in Te6edf3message Tff7b72&Tff7b72&
Tb4b4b4(Te6edf3ROOM_READER_CONNECTION_PHRASE Tff7b72in Te6edf3message Tff7b72|Tff7b72| Te6edf3ROOM_WRITER_CONNECTION_PHRASE Tff7b72in Te6edf3messageTb4b4b4)
Tff7b72fun Td2a8ffisDbPoolAcquireTimeoutExceptionTb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3generateSequenceTff7b72<Te6edf3ThrowableTff7b72>Tb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cause Tb4b4b4}
Tb4b4b4.Te6edf3any Tb4b4b4{ Te6edf3throwable Tff7b72-Tff7b72>
Tff7b72val Te6edf3msg Tff7b72= Te6edf3throwableTb4b4b4.Te6edf3messageTff7b72?.Te6edf3lowercaseTb4b4b4(Tb4b4b4) Tff7b72?: Tff7b72returnTf0883e@any Tff7b72false
Te6edf3isRoomPoolAcquireTimeoutMessageTb4b4b4(Te6edf3msgTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Returns true if a database exists for the given device address. Android Room stores DB files without an
* extension; JVM/iOS append `.db`. We check both to stay platform-agnostic.
*/
Tff7b72override Tff7b72fun Td2a8ffhasDatabaseForTb4b4b4(Te6edf3addressTb4b4b4: Tffa657String?Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3addressTb4b4b4.Te6edf3isNullOrBlankTb4b4b4(Tb4b4b4) Tff7b72|Tff7b72| Te6edf3address Tff7b72=Tff7b72= Ta5d6ff"Ta5d6ffnTa5d6ff"Tb4b4b4) Tff7b72return Tff7b72false
Tff7b72val Te6edf3dbName Tff7b72= Te6edf3buildDbNameTb4b4b4(Te6edf3addressTb4b4b4)
Tff7b72return Te6edf3dbFileExistsTb4b4b4(Te6edf3dbNameTb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffdbFileExistsTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3dir Tff7b72= Te6edf3getDatabaseDirectoryTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3fs Tff7b72= Te6edf3getFileSystemTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3fsTb4b4b4.Te6edf3existsTb4b4b4(Te6edf3dirTb4b4b4.Te6edf3resolveTb4b4b4(Te6edf3dbNameTb4b4b4)Tb4b4b4) Tff7b72|Tff7b72| Te6edf3fsTb4b4b4.Te6edf3existsTb4b4b4(Te6edf3dirTb4b4b4.Te6edf3resolveTb4b4b4(Ta5d6ff"Tffd700$Te6edf3dbNameTa5d6ff.dbTa5d6ff"Tb4b4b4)Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffdbFileMetadataMillisTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Long? Tb4b4b4{
Tff7b72val Te6edf3dir Tff7b72= Te6edf3getDatabaseDirectoryTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3fs Tff7b72= Te6edf3getFileSystemTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3fsTb4b4b4.Te6edf3metadataOrNullTb4b4b4(Te6edf3dirTb4b4b4.Te6edf3resolveTb4b4b4(Te6edf3dbNameTb4b4b4)Tb4b4b4)Tff7b72?.Te6edf3lastModifiedAtMillis
Tff7b72?: Te6edf3fsTb4b4b4.Te6edf3metadataOrNullTb4b4b4(Te6edf3dirTb4b4b4.Te6edf3resolveTb4b4b4(Ta5d6ff"Tffd700$Te6edf3dbNameTa5d6ff.dbTa5d6ff"Tb4b4b4)Tb4b4b4)Tff7b72?.Te6edf3lastModifiedAtMillis
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffmarkLastUsedTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Te6edf3launchManagerWork Tb4b4b4{ Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTff7b72[Te6edf3lastUsedKeyTb4b4b4(Te6edf3dbNameTb4b4b4)Tff7b72] Tff7b72= Te6edf3nowMillis Tb4b4b4} Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8fflastUsedTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4)Tb4b4b4: Tffa657Long Tb4b4b4{
Tff7b72val Te6edf3key Tff7b72= Te6edf3lastUsedKeyTb4b4b4(Te6edf3dbNameTb4b4b4)
Tff7b72val Te6edf3v Tff7b72= Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tff7b72[Te6edf3keyTff7b72] Tff7b72?: T79c0ff0L
Tff7b72return Tff7b72if Tb4b4b4(Te6edf3v Tff7b72=Tff7b72= T79c0ff0LTb4b4b4) Tb4b4b4{
Te6edf3dbFileMetadataMillisTb4b4b4(Te6edf3dbNameTb4b4b4) Tff7b72?: T79c0ff0L
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3v
Tb4b4b4}
Tb4b4b4}
Tff7b72protected Tff7b72open Tff7b72fun Td2a8fflistExistingDbNamesTb4b4b4(Tb4b4b4)Tb4b4b4: Te6edf3ListTff7b72<Tffa657StringTff7b72> Tb4b4b4{
Tff7b72val Te6edf3dir Tff7b72= Te6edf3getDatabaseDirectoryTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3fs Tff7b72= Te6edf3getFileSystemTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3fsTb4b4b4.Te6edf3existsTb4b4b4(Te6edf3dirTb4b4b4)Tb4b4b4) Tff7b72return Te6edf3emptyListTb4b4b4(Tb4b4b4)
Tff7b72return Te6edf3fsTb4b4b4.Te6edf3listTb4b4b4(Te6edf3dirTb4b4b4)
Tb4b4b4.Te6edf3asSequenceTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3name Tb4b4b4}
Tb4b4b4.Te6edf3filter Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3startsWithTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DB_PREFIXTb4b4b4) Tb4b4b4}
T8b949e// Skip Room-internal sidecar files (-wal/-shm/-journal) and lock files so each DB appears exactly once.
Tb4b4b4.Te6edf3filterNot Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3endsWithTb4b4b4(Ta5d6ff"Ta5d6ff-walTa5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Tffa657itTb4b4b4.Te6edf3endsWithTb4b4b4(Ta5d6ff"Ta5d6ff-shmTa5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Tffa657itTb4b4b4.Te6edf3endsWithTb4b4b4(Ta5d6ff"Ta5d6ff-journalTa5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Tffa657itTb4b4b4.Te6edf3endsWithTb4b4b4(Ta5d6ff"Ta5d6ff.lckTa5d6ff"Tb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3removeSuffixTb4b4b4(Ta5d6ff"Ta5d6ff.dbTa5d6ff"Tb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3distinctTb4b4b4(Tb4b4b4)
Tb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffenforceCacheLimitTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Deferred enforcement can wait behind a later switch. Resolve the protected name under the same mutex
T8b949e// that publishes currentDbName so the active database at execution time can never become an LRU victim.
Tff7b72val Te6edf3activeDbName Tff7b72= Te6edf3currentDbName
Tff7b72val Te6edf3limit Tff7b72= Te6edf3getCurrentCacheLimitTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3all Tff7b72= Te6edf3listExistingDbNamesTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3pendingRouteNames Tff7b72= Te6edf3pendingRouteDbNamesTb4b4b4(Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tb4b4b4)
Tff7b72val Te6edf3detachedDbNames Tff7b72= Te6edf3detachedDatabasesTb4b4b4.Te6edf3mapToTb4b4b4(Te6edf3mutableSetOfTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3dbName Tb4b4b4}
T8b949e// Only enforce the limit over device-specific DBs. A detached pool is still live for a consumer of an
T8b949e// earlier currentDb emission, so its files must remain protected until orderly shutdown.
Tff7b72val Te6edf3deviceDbs Tff7b72=
Te6edf3allTb4b4b4.Te6edf3filterNot Tb4b4b4{
Tffa657it Tff7b72in Te6edf3logicallyRetired Tff7b72|Tff7b72|
Tffa657it Tff7b72in Te6edf3detachedDbNames Tff7b72|Tff7b72|
Tffa657it Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3LEGACY_DB_NAME Tff7b72|Tff7b72|
Tffa657it Tff7b72=Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAME
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3deviceDbsTb4b4b4.Te6edf3size Tff7b72<Tff7b72= Te6edf3limitTb4b4b4) Tff7b72returnTf0883e@withLock
Tff7b72val Te6edf3usageSnapshot Tff7b72= Te6edf3deviceDbsTb4b4b4.Te6edf3associateWith Tb4b4b4{ Te6edf3lastUsedTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
T8b949e// A pending destination can be the only merged copy and its merge marker is the proof needed to repair
T8b949e// the route after a crash. Keep both route endpoints until address-scoped recovery finalizes or clears it.
Tff7b72val Te6edf3victims Tff7b72=
Te6edf3selectEvictionVictimsTb4b4b4(
Te6edf3dbNames Tff7b72= Te6edf3deviceDbsTb4b4b4,
Te6edf3activeDbName Tff7b72= Te6edf3activeDbNameTb4b4b4,
Te6edf3limit Tff7b72= Te6edf3limitTb4b4b4,
Te6edf3lastUsedMsByDb Tff7b72= Te6edf3usageSnapshotTb4b4b4,
Te6edf3protectedDbNames Tff7b72= Te6edf3pendingRouteNamesTb4b4b4,
Tb4b4b4)
Tff7b72val Te6edf3evictableVictims Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3victimsTb4b4b4.Te6edf3filter Tb4b4b4{ Te6edf3name Tff7b72-Tff7b72>
Tff7b72val Te6edf3cached Tff7b72= Te6edf3dbCacheTff7b72[Te6edf3nameTff7b72]
Tff7b72if Tb4b4b4(Te6edf3cached Tff7b72!Tff7b72= Tff7b72null Tff7b72&Tff7b72& Te6edf3hasActiveDatabaseAccessLockedTb4b4b4(Te6edf3cachedTb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3deferredEvictionsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3cachedTb4b4b4)
Tff7b72false
Tb4b4b4} Tff7b72else Tb4b4b4{
Tff7b72true
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3evictableVictimsTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3name Tff7b72-Tff7b72>
Tff7b72try Tb4b4b4{
Te6edf3closeCachedDatabaseTb4b4b4(Te6edf3nameTb4b4b4)
Te6edf3deleteDatabaseFilesTb4b4b4(Te6edf3nameTb4b4b4)
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3lastUsedKeyTb4b4b4(Te6edf3nameTb4b4b4)Tb4b4b4) Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffEvicted cached DB Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3nameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3cancellationTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3cancellation
Tb4b4b4} Tff7b72catch Tb4b4b4(Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4) Te6edf3failureTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to evict database Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3nameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffcleanupLegacyDbIfNeededTb4b4b4(Te6edf3activeDbNameTb4b4b4: Tffa657StringTb4b4b4) Tff7b72= Te6edf3withManagerOperation Tb4b4b4{
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3cleaned Tff7b72= Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tff7b72[Te6edf3legacyCleanedKeyTff7b72] Tff7b72?: Tff7b72false
Tff7b72if Tb4b4b4(Te6edf3cleanedTb4b4b4) Tff7b72returnTf0883e@withLock
Tff7b72val Te6edf3legacy Tff7b72= Te6edf3DatabaseConstantsTb4b4b4.Te6edf3LEGACY_DB_NAME
Tff7b72if Tb4b4b4(Te6edf3legacy Tff7b72=Tff7b72= Te6edf3activeDbNameTb4b4b4) Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTff7b72[Te6edf3legacyCleanedKeyTff7b72] Tff7b72= Tff7b72true Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3dbFileExistsTb4b4b4(Te6edf3legacyTb4b4b4)Tb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Te6edf3closeCachedDatabaseTb4b4b4(Te6edf3legacyTb4b4b4)
Te6edf3deleteDatabaseFilesTb4b4b4(Te6edf3legacyTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3cancellationTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3cancellation
Tb4b4b4} Tff7b72catch Tb4b4b4(Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4) Te6edf3failureTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to delete legacy database Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3legacyTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffDeleted legacy DB Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3legacyTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Te6edf3datastoreTb4b4b4.Te6edf3edit Tb4b4b4{ Tffa657itTff7b72[Te6edf3legacyCleanedKeyTff7b72] Tff7b72= Tff7b72true Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72fun Td2a8ffscheduleSearchIndexBackfillTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4, Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4, Te6edf3shouldDelayBackfillTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Te6edf3backfillJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3backfillJob Tff7b72=
Te6edf3launchManagerWorkTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3shouldDelayBackfillTb4b4b4) Te6edf3delayTb4b4b4(Te6edf3BACKFILL_COLD_START_DELAY_MSTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3_currentDbTb4b4b4.Te6edf3value Tff7b72!Tff7b72=Tff7b72= Te6edf3dbTb4b4b4) Tff7b72returnTf0883e@launchManagerWork
Te6edf3backfillSearchIndexIfNeededTb4b4b4(Te6edf3dbTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to backfill search index for Tffd700${Te6edf3anonymizeDbNameTb4b4b4(Te6edf3dbNameTb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Backfills [Packet.messageText] for existing text-message packets that predate the FTS5 schema, then rebuilds the
* FTS index so search covers historical messages. The text is decoded in Kotlin from each packet's payload (see
* [PacketDao.backfillMessageTexts]); it cannot be read in SQL because the message body is stored as serialized
* `bytes`, not a `text` JSON field.
*/
Tff7b72internal Tff7b72suspend Tff7b72fun Td2a8ffbackfillSearchIndexIfNeededTb4b4b4(Te6edf3scheduledDbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3admission Tff7b72= Te6edf3beginWriteTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3admittedDb Tff7b72= Te6edf3admissionTb4b4b4.Te6edf3database
Tff7b72try Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3admittedDb Tff7b72!Tff7b72=Tff7b72= Te6edf3scheduledDbTb4b4b4) Tff7b72return
Te6edf3runCancellableDbBlockTb4b4b4(Te6edf3admittedDbTb4b4b4) Tb4b4b4{ Te6edf3performSearchIndexBackfillTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{ Te6edf3endWriteTb4b4b4(Te6edf3admittedDbTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Performs count, message-text backfill, and FTS rebuild while caller holds one writer admission. */
Tff7b72protected Tff7b72open Tff7b72suspend Tff7b72fun Td2a8ffperformSearchIndexBackfillTb4b4b4(Te6edf3dbTb4b4b4: Te6edf3MeshtasticDatabaseTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3dbTb4b4b4.Te6edf3packetDaoTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3countPacketsNeedingBackfillTb4b4b4(Tb4b4b4) Tff7b72=Tff7b72= T79c0ff0Tb4b4b4) Tff7b72return
Tff7b72val Te6edf3count Tff7b72= Te6edf3dbTb4b4b4.Te6edf3packetDaoTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3backfillMessageTextsTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3count Tff7b72> T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffBackfilled Tffd700$Te6edf3countTa5d6ff messages for FTS search indexTa5d6ff" Tb4b4b4}
Te6edf3dbTb4b4b4.Te6edf3packetDaoTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3rebuildFtsIndexTb4b4b4(Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffFTS search index rebuild completeTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/** Platform file removal seam used by orderly retirement and test fixtures. */
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffdeleteDatabaseFilesTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4) Tff7b72= Te6edf3deleteDatabaseTb4b4b4(Te6edf3dbNameTb4b4b4)
T8b949e/**
* Establishes an orderly shutdown boundary: rejects new work, bounds manager-job cancellation, admitted
* manager-operation draining, and admitted-writer draining, waits for the last serialized switch/association to
* finalize, then closes every manager-owned Room instance. If a cancelled child, admitted operation, writer, or
* pool close cannot finish successfully, ownership is retained and physical cleanup is skipped so a later [close]
* call can retry without losing track of live resources. Retried attempts may call Room's idempotent `close()`
* again for pools that completed during an earlier partial attempt.
*/
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffCyclomaticComplexMethodTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffLongMethodTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72suspend Tff7b72fun Td2a8ffcloseTb4b4b4(Tb4b4b4) Tff7b72= Te6edf3withContextTb4b4b4(Te6edf3NonCancellableTb4b4b4) Tb4b4b4{
Te6edf3closeMutexTb4b4b4.Te6edf3withLock Tf0883ecloseAttempt@Tb4b4b4{
Tff7b72var Te6edf3operationsDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
Tff7b72var Te6edf3writerDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
Tff7b72var Te6edf3jobsDrainTb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tff7b72? Tff7b72= Tff7b72null
Tff7b72var Te6edf3closedSuccessfully Tff7b72= Tff7b72false
Tff7b72val Te6edf3shouldClose Tff7b72=
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72when Tb4b4b4(Te6edf3lifecycleStateTb4b4b4) Tb4b4b4{
Te6edf3LifecycleStateTb4b4b4.Te6edf3CLOSED Tff7b72-Tff7b72> Tff7b72false
Te6edf3LifecycleStateTb4b4b4.Te6edf3OPENTb4b4b4,
Te6edf3LifecycleStateTb4b4b4.Te6edf3CLOSINGTb4b4b4,
Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3lifecycleState Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3CLOSING
Te6edf3operationsDrain Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3activeManagerOperationsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3managerOperationDrain Tff7b72= Tffa657it Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3managerOperationDrain Tff7b72= Tff7b72null
Tff7b72null
Tb4b4b4}
Te6edf3writerDrain Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3activeWritersTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3shutdownWriterDrain Tff7b72= Tffa657it Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3shutdownWriterDrain Tff7b72= Tff7b72null
Tff7b72null
Tb4b4b4}
Tff7b72true
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3shouldCloseTb4b4b4) Tff7b72returnTf0883e@closeAttempt
Tff7b72val Te6edf3managerJobs Tff7b72=
Te6edf3synchronizedTb4b4b4(Te6edf3managerJobLockTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3jobs Tff7b72= Te6edf3activeManagerJobsTb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4)
Te6edf3jobsDrain Tff7b72=
Tff7b72if Tb4b4b4(Te6edf3jobsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3CompletableDeferredTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3managerJobDrain Tff7b72= Tffa657it Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3managerJobDrain Tff7b72= Tff7b72null
Tff7b72null
Tb4b4b4}
Te6edf3jobs
Tb4b4b4}
Tff7b72try Tb4b4b4{
T8b949e// Manager-owned cleanup jobs hold manager-operation tokens. Publish cancellation before waiting for
T8b949e// those operations so cancellable I/O can unwind instead of forcing every close attempt to time
T8b949e// out.
Te6edf3managerJobsTb4b4b4.Te6edf3forEach Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cancelTb4b4b4(Tb4b4b4) Tb4b4b4}
T8b949e// Bound-wait for already-admitted manager operations before acquiring [mutex]. New work is rejected
T8b949e// while CLOSING, and a timed-out attempt leaves ownership intact so a later close() can retry.
Tff7b72val Te6edf3operationsDrained Tff7b72=
Te6edf3operationsDrainTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3drain Tff7b72-Tff7b72>
Tff7b72val Te6edf3completedBeforeTimeout Tff7b72=
Te6edf3withContextTb4b4b4(Te6edf3DispatchersTb4b4b4.Te6edf3DefaultTb4b4b4) Tb4b4b4{
Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3WRITER_DRAIN_TIMEOUT_MSTb4b4b4) Tb4b4b4{
Te6edf3drainTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4)
Tff7b72true
Tb4b4b4} Tff7b72?: Tff7b72false
Tb4b4b4}
Te6edf3completedBeforeTimeout Tff7b72|Tff7b72| Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3activeManagerOperationsTb4b4b4.Te6edf3isEmptyTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72?: Tff7b72true
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3operationsDrainedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffDatabase shutdown timed out waiting for in-flight manager operations; retaining owned pools Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6fffor a later close retryTa5d6ff"
Tb4b4b4}
Tff7b72returnTf0883e@closeAttempt
Tb4b4b4}
Tff7b72val Te6edf3managerJobsStopped Tff7b72=
Te6edf3jobsDrainTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3WRITER_DRAIN_TIMEOUT_MSTb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72!Tff7b72= Tff7b72null Tb4b4b4} Tff7b72?: Tff7b72true
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3managerJobsStoppedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffTimed out stopping database manager jobs; retaining owned pools for a later close retryTa5d6ff"
Tb4b4b4}
Tff7b72returnTf0883e@closeAttempt
Tb4b4b4}
Tff7b72val Te6edf3writersDrained Tff7b72=
Te6edf3writerDrainTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3WRITER_DRAIN_TIMEOUT_MSTb4b4b4) Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72!Tff7b72= Tff7b72null Tb4b4b4} Tff7b72?: Tff7b72true
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3writersDrainedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffDatabase shutdown could not prove every writer stopped; retaining owned pools for a later Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffclose retryTa5d6ff"
Tb4b4b4}
Tff7b72returnTf0883e@closeAttempt
Tb4b4b4}
T8b949e// Snapshot ownership without transferring it yet. A failed pool close must leave every instance and
T8b949e// retirement intent reachable by a later close() attempt.
Tff7b72val Te6edf3snapshot Tff7b72=
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3databases Tff7b72= Te6edf3mutableListOfTff7b72<Te6edf3ShutdownDatabaseTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72fun Td2a8ffaddDistinctTb4b4b4(Te6edf3dbNameTb4b4b4: Tffa657StringTb4b4b4, Te6edf3databaseTb4b4b4: Te6edf3MeshtasticDatabase?Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3database Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tff7b72return
Tff7b72val Te6edf3existing Tff7b72= Te6edf3databasesTb4b4b4.Te6edf3firstOrNull Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3database Tff7b72=Tff7b72=Tff7b72= Te6edf3database Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3existing Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3databasesTb4b4b4.Te6edf3addTb4b4b4(Te6edf3ShutdownDatabaseTb4b4b4(Te6edf3databaseTb4b4b4, Te6edf3mutableSetOfTb4b4b4(Te6edf3dbNameTb4b4b4)Tb4b4b4)Tb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3existingTb4b4b4.Te6edf3dbNamesTb4b4b4.Te6edf3addTb4b4b4(Te6edf3dbNameTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Te6edf3dbCacheTb4b4b4.Te6edf3forEach Tb4b4b4{ Tb4b4b4(Te6edf3dbNameTb4b4b4, Te6edf3databaseTb4b4b4) Tff7b72-Tff7b72> Te6edf3addDistinctTb4b4b4(Te6edf3dbNameTb4b4b4, Te6edf3databaseTb4b4b4) Tb4b4b4}
Te6edf3detachedDatabasesTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3addDistinctTb4b4b4(Tffa657itTb4b4b4.Te6edf3dbNameTb4b4b4, Tffa657itTb4b4b4.Te6edf3databaseTb4b4b4) Tb4b4b4}
Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{
Te6edf3addDistinctTb4b4b4(Te6edf3DatabaseConstantsTb4b4b4.Te6edf3DEFAULT_DB_NAMETb4b4b4, Te6edf3initializedDefaultDbTb4b4b4)
Te6edf3addDistinctTb4b4b4(Te6edf3currentDbNameTb4b4b4, Te6edf3currentDbStateTff7b72?.Te6edf3valueTb4b4b4)
Tb4b4b4}
Tff7b72val Te6edf3persistedRetiredNames Tff7b72=
Tff7b72try Tb4b4b4{
Te6edf3datastoreTb4b4b4.Te6edf3dataTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tff7b72[Te6edf3retiredDbNamesKeyTff7b72]Tb4b4b4.Te6edf3orEmptyTb4b4b4(Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3failureTb4b4b4: Te6edf3ThrowableTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Ta5d6ff"Ta5d6ffFailed to read persisted database retirements during shutdownTa5d6ff"
Tb4b4b4}
Te6edf3emptySetTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3ShutdownSnapshotTb4b4b4(Te6edf3databasesTb4b4b4, Tb4b4b4(Te6edf3logicallyRetired Tff7b72+ Te6edf3persistedRetiredNamesTb4b4b4)Tb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4)Tb4b4b4)
Tb4b4b4}
T8b949e// All tracked work has stopped. The remaining scope-owned collector does not touch Room; cancel it
T8b949e// before closing pools, but keep ownership maps intact until every close succeeds.
Te6edf3managerScopeTb4b4b4.Te6edf3coroutineContextTff7b72[Te6edf3JobTff7b72]Tff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3backfillJob Tff7b72= Tff7b72null
Tff7b72val Te6edf3failedCloseNames Tff7b72= Te6edf3mutableSetOfTff7b72<Tffa657StringTff7b72>Tb4b4b4(Tb4b4b4)
Te6edf3snapshotTb4b4b4.Te6edf3databasesTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3owned Tff7b72-Tff7b72>
Te6edf3runCatching Tb4b4b4{ Te6edf3closeDatabaseTb4b4b4(Te6edf3ownedTb4b4b4.Te6edf3databaseTb4b4b4) Tb4b4b4}
Tb4b4b4.Te6edf3onFailure Tb4b4b4{
Te6edf3failedCloseNames Tff7b72+Tff7b72= Te6edf3ownedTb4b4b4.Te6edf3dbNames
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to close database during shutdownTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3failedCloseNamesTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{
Ta5d6ff"Ta5d6ffDatabase shutdown retained Tffd700${Te6edf3failedCloseNamesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff pool name(s) after close failure; Ta5d6ff" Tff7b72+
Ta5d6ff"Ta5d6ffa later close() will retryTa5d6ff"
Tb4b4b4}
Tff7b72returnTf0883e@closeAttempt
Tb4b4b4}
Te6edf3snapshotTb4b4b4.Te6edf3retiredNamesTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3physicallyRetireDatabaseTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
T8b949e// Only now transfer ownership and publish the terminal state. No admitted work can add another pool
T8b949e// because the manager has remained CLOSING throughout this attempt.
Te6edf3mutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3dbCacheTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3detachedDatabasesTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3logicallyRetiredTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3synchronizedTb4b4b4(Te6edf3initializationLockTb4b4b4) Tb4b4b4{
Te6edf3initializedDefaultDb Tff7b72= Tff7b72null
Te6edf3currentDbState Tff7b72= Tff7b72null
Tb4b4b4}
Tb4b4b4}
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3lifecycleState Tff7b72= Te6edf3LifecycleStateTb4b4b4.Te6edf3CLOSED
Te6edf3deferredEvictionsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3drainWaitersTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3activeWritersTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3activeReadersTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3poolLanesTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3unrecoverablePoolsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3activeManagerOperationsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3writerGateTff7b72?.Te6edf3completeExceptionallyTb4b4b4(Te6edf3IllegalStateExceptionTb4b4b4(Ta5d6ff"Ta5d6ffDatabaseManager is closing or closedTa5d6ff"Tb4b4b4)Tb4b4b4)
Te6edf3writerGate Tff7b72= Tff7b72null
Tb4b4b4}
Te6edf3closedSuccessfully Tff7b72= Tff7b72true
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3synchronizedTb4b4b4(Te6edf3managerJobLockTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3managerJobDrain Tff7b72=Tff7b72=Tff7b72= Te6edf3jobsDrainTb4b4b4) Te6edf3managerJobDrain Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3closedSuccessfullyTb4b4b4) Te6edf3activeManagerJobsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3writerTrackerMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3managerOperationDrain Tff7b72=Tff7b72=Tff7b72= Te6edf3operationsDrainTb4b4b4) Te6edf3managerOperationDrain Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Te6edf3shutdownWriterDrain Tff7b72=Tff7b72=Tff7b72= Te6edf3writerDrainTb4b4b4) Te6edf3shutdownWriterDrain Tff7b72= Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* A bounded one-shot database operation did not finish within its deadline and was abandoned, never replayed.
*
* Room's connection pool retries a failed acquisition indefinitely instead of throwing, so a leaked permit surfaces as
* a callback that never returns. This is the deadline that turns that silent hang into a failure the caller can react
* to; the active pool is reopened so subsequent operations proceed.
*/
Tff7b72class T56d364DatabaseOperationTimeoutExceptionTb4b4b4(Te6edf3messageTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4: Te6edf3IllegalStateExceptionTb4b4b4(Te6edf3messageTb4b4b4)
Served by rngit 1.5.0 - Generated in 0.36s